I am not seeing records in the consuming segments...
# general
a
I am not seeing records in the consuming segments for 1 shard of kinesis workflowEvents__1__444__20230123T2054Z
Copy code
{
  "segment.creation.time": "1674507284747",
  "segment.flush.threshold.size": "5000000",
  "segment.name": "workflowEvents__1__444__20230123T2054Z",
  "segment.realtime.numReplicas": "3",
  "segment.realtime.startOffset": "{\"shardId-000000000001\":\"49632676210804516725069199056584947625655896387241377810\"}",
  "segment.realtime.status": "IN_PROGRESS",
  "segment.table.name": "workflowEvents",
  "segment.type": "REALTIME"
}
My kinesis has 2 shards 1 of the segments for shard 0 are up to date but segement for shard 1 are delayed
Previous shard workflowEvents__1__443__20230122T2054Z
Copy code
{
  "segment.crc": "2300299016",
  "segment.creation.time": "1674420871153",
  "segment.end.time": "1674507223224",
  "segment.flush.threshold.size": "5000000",
  "segment.index.version": "v3",
  "segment.name": "workflowEvents__1__443__20230122T2054Z",
  "segment.realtime.download.url": "<s3://ctct-transient-prod-us-east-1-eigi-datalake/pinot/controller-data/workflowEvents/workflowEvents__1__443__20230122T2054Z>",
  "segment.realtime.endOffset": "{\"shardId-000000000001\":\"49632676210804516725069199056584947625655896387241377810\"}",
  "segment.realtime.numReplicas": "3",
  "segment.realtime.startOffset": "{\"shardId-000000000001\":\"49632676210804516725069195830507160618850891209447047186\"}",
  "segment.realtime.status": "DONE",
  "segment.start.time": "1669563719851",
  "segment.table.name": "workflowEvents",
  "segment.time.unit": "MILLISECONDS",
  "segment.total.docs": "119507",
  "segment.type": "REALTIME"
}
(edited)
x
from the metadata, this segment is already sealed.
Copy code
"segment.realtime.status": "DONE",
is there a new
workflowEvents__1__444__2023012XTXXXXZ
segment?
a
Yes I can only see records in this shared but none in the consuming shards and I have posted records for an account which is tied to that shard after the closing Time
x
the controller should fix this problem during the validation phase. cc: @Neha Pawar
a
workflowEvents__1__444__20230123T2054Z is the new consuming segment
@Xiangfu 1234
x
ok, so I guess this is working now ?
a
No it is not the new consuming segment does not have any records and I have sent records to the shard with timestamp higher than start time..that is the issue
@Xiangfu 1234 the records are only being queryable after the segment is closed but only for segment with partition 1 segment with partition 0 are showing immediately
Copy code
{
  "REALTIME": {
    "tableName": "workflowEvents_REALTIME",
    "tableType": "REALTIME",
    "segmentsConfig": {
      "timeType": "MILLISECONDS",
      "schemaName": "workflowEvents",
      "retentionTimeUnit": "DAYS",
      "retentionTimeValue": "1826",
      "allowNullTimeValue": false,
      "replicasPerPartition": "3",
      "segmentPushType": "APPEND",
      "timeColumnName": "eventTimestamp"
    },
    "tenants": {
      "broker": "DefaultTenant",
      "server": "DefaultTenant"
    },
    "tableIndexConfig": {
      "streamConfigs": {
        "streamType": "kinesis",
        "stream.kinesis.topic.name": "prod-rel-cdp-dl-workflow-metrics-stream",
        "region": "us-east-1",
        "shardIteratorType": "LATEST",
        "stream.kinesis.consumer.type": "lowlevel",
        "stream.kinesis.fetch.timeout.millis": "30000",
        "stream.kinesis.decoder.class.name": "org.apache.pinot.plugin.stream.kafka.KafkaJSONMessageDecoder",
        "stream.kinesis.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kinesis.KinesisConsumerFactory",
        "realtime.segment.flush.threshold.size": "5000000",
        "realtime.segment.flush.threshold.time": "1d"
      },
      "rangeIndexVersion": 1,
      "autoGeneratedInvertedIndex": false,
      "createInvertedIndexDuringSegmentGeneration": false,
      "loadMode": "MMAP",
      "enableDefaultStarTree": false,
      "enableDynamicStarTreeCreation": false,
      "aggregateMetrics": false,
      "nullHandlingEnabled": false
    },
    "metadata": {
      "customConfigs": {}
    },
    "routing": {
      "instanceSelectorType": "strictReplicaGroup"
    },
    "upsertConfig": {
      "mode": "FULL",
      "hashFunction": "NONE"
    },
    "isDimTable": false
  }
}
@Kartik Khare @Xiangfu 1234 @Neha Pawar any ideas why I am only seeing records for segments
workflowEvents__1__*
only after flush ?
x
Hmm, can you try to disable and enable the table? Also check any error in the server hosting the consuming segment but without any data
a
Only exception i see is this
Copy code
The request rate for the stream is too high                                                                                                                                  │
│ shaded.software.amazon.awssdk.services.kinesis.model.ProvisionedThroughputExceededException: Rate exceeded for Shard -  │
│     at shaded.software.amazon.awssdk.core.internal.http.CombinedResponseHandler.handleErrorResponse(CombinedResponseHandler.java:123) ~[pinot-all-0.9.1-jar-with-dependencie │
│     at shaded.software.amazon.awssdk.core.internal.http.CombinedResponseHandler.handleResponse(CombinedResponseHandler.java:79) ~[pinot-all-0.9.1-jar-with-dependencies.jar: │
│     at
x
So it’s AWS side rate limit?
a
yes kinesis
I also see these log messages
Copy code
Processed requestId=55106106,table=workflowEvents_REALTIME,segments(queried/processed/matched/consuming)=444/115/0/-1,schedulerWaitMs=0,reqDeserMs=0,totalExecMs=3,resSerMs= │
│
But I see the same on the other servers as well
I also see these logs which means it is consuming but it is not accessible
Copy code
2023/01/24 20:39:47.982 INFO [LLRealtimeSegmentDataManager_workflowEvents__1__444__20230123T2054Z] [workflowEvents__1__444__20230123T2054Z] Consumed 9 events from (rate:0.1485026/s), currentOffset={"shardId-000000000001":"49632676210804516725069202648418405653588441169257299986"}, numRowsConsumedSoFar=213200, numRowsIndexedSoFar=213200
x
Hmm so I guess it’s just slow?
Can you increase the kinesis side limit?
a
but it has been almost 24 hours and i am not seeing a single record in workflowEvents__1__444__20230123T2054Z ..I see many of the above logs =
x
Right, likely to be request exceed the limit, so it’s never proceed
For consumed msgs, do you mean those messages are consumed but no message show up in Pinot?
If that is the case, then the only thing I can think of is that the message itself changed and Pinot cannot decode it correctly
a
Yes
x
I would anyway go ahead to increase the AWS kinesis rate limit
Then maybe create a console kinesis consumer to consume message and see if there is any changes
a
I just verified the pinot ,,I am able to query the records now after 1 day which means that the records are visible once they are flushed
To reiterate the problem I am able to query records from consuming segments workflowEvents__0 but not workflowEvents__1 and the records are the same
@Xiangfu 1234 could this be related to this problem https://apache-pinot.slack.com/archives/CDRCA57FC/p1674515086070049
@Xiangfu 1234 I rebalanced the cluster yesterday but it still shows Bad . Could this be the result of consuming segment being out of status
Copy code
IDEALSTATES
"workflowEvents__0__84__20220205T0439Z": {
      "Server_cdp-dl-pinot-k8s-server-2.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local_8098": "ONLINE",
      "Server_cdp-dl-pinot-k8s-server-3.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local_8098": "ONLINE",
      "Server_cdp-dl-pinot-k8s-server-4.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local_8098": "ONLINE"
    },

EXTERNALVIEW of same 
"workflowEvents__0__80__20220201T0438Z": {
      "Server_cdp-dl-pinot-k8s-server-0.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local_8098": "ONLINE",
      "Server_cdp-dl-pinot-k8s-server-2.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local_8098": "ONLINE",
      "Server_cdp-dl-pinot-k8s-server-3.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local_8098": "ONLINE",
      "Server_cdp-dl-pinot-k8s-server-4.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local_8098": "ONLINE"
    },
how to fix this ?
n
can you share the whole ideal state and external view? this looks like it’s for different segments. also check your server logs, the external view has more servers than the replication and the server logs should tell you why that’s happening
a
I restarted the server and now I dont see 4 replicas but I do see a few segments are stuck in 2/3 state
Screen Shot 2023-01-25 at 3.41.29 PM.png
@Neha Pawar unfortunately i copied the wrong segment but fortunately that problem is gone after a restart