hi - we have two pinot tables both consuming from ...
# troubleshooting
s
hi - we have two pinot tables both consuming from the same Kafka topic. Both are using the low level consumer. One is a hybrid table and one is a realtime only table. We have an issue where one record is missing from the realtime table but is present in the hybrid table. We have looked in the logs and can see no Warn or Error messages at the time the record was lost. The only log of interest is that we have an idle consumer at that time and the stream is recreated. Are there any known scenarios in which message loss is possible?
d
I'll top up here with some additional findings
out schema looks like this
Copy code
"schemaName": "audit_event",

    "primaryKeyColumns": [

      "message_key"

    ],

    "dimensionFieldSpecs": [

      {

        "name": "message_key",

        "dataType": "STRING"

      },

      ...

    ],

    "dateTimeFieldSpecs": [

      {

        "name": "event_timestamp",

        "dataType": "LONG",

        "format": "1:MILLISECONDS:EPOCH",

        "granularity": "1:DAYS"

      },

      {

        "name": "audit_timestamp",

        "dataType": "LONG",

        "format": "1:MILLISECONDS:EPOCH",

        "granularity": "1:DAYS"

      }

    ]

  }
the table definition is
Copy code
{
  "tableName": "audit_event",
  "tableType": "REALTIME",
  "routing": {
    "segmentPrunerTypes": [
      "time"
    ]
  },
  "task": {
    "taskTypeConfigsMap": {
      "RealtimeToOfflineSegmentsTask": {
        "bucketTimePeriod": "1d",
        "bufferTimePeriod": "2d",
        "mergeType": "dedup",
        "maxNumRecordsPerSegment": "10000000"
      }
    }
  },
  "segmentsConfig": {
    "schemaName": "audit_event ",
    "timeColumnName": "event_timestamp",
    "timeType": "MILLISECONDS",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": 5,
    "replication": 2,
    "replicasPerPartition": 2,
    "completionMode": "DOWNLOAD",
    "peerSegmentDownloadScheme": "https"
  },
  "tableIndexConfig": {
    "createInvertedIndexDuringSegmentGeneration": false,
    "enableDefaultStarTree": false,
    "enableDynamicStarTreeCreation": false,
    "loadMode": "MMAP",
    "columnMinMaxValueGeneratorMode": "NONE",
    "nullHandlingEnabled": true,
    "aggregateMetrics": false
  },
  "tenants": {
    "broker": "DefaultTenant",
    "server": "DefaultTenant",
    "tagOverrideConfig": {}
  },
  "ingestionConfig": {
    "streamIngestionConfig": {
      "streamConfigMaps": [
        {
          "streamType": "kafka",
          "stream.kafka.consumer.type": "LowLevel",
          "stream.kafka.topic.name": "audit-event",
          "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder",
          "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory",
          "stream.kafka.decoder.prop.schema.registry.rest.url": "${KAFKA_SCHEMA_REGISTRY_URL}",
          "stream.kafka.broker.list": "${KAFKA_BOOTSTRAP_SERVERS}",
          "stream.kafka.consumer.prop.auto.offset.reset": "smallest",
          "sasl.jaas.config": "${SASL_JAAS_CONFIG}",
          "security.protocol": "SASL_SSL",
          "sasl.mechanism": "PLAIN",
          "endpoint.identification.algorithm": "https",
          "realtime.segment.flush.threshold.rows": "0",
          "realtime.segment.flush.autotune.initialRows": "100000",
          "realtime.segment.flush.threshold.time": "12h",
          "realtime.segment.flush.threshold.segment.size": "250M"
        }
      ]
    }
  },
  "metadata": {
    "customConfigs": {
      "team": "my team"
    }
  }
}
this is a small cluster with 2 servers only
we have reconciled the data in pinot with our data lake and we have found that the missing records are always the first record in the segment. One of these records has event timestamp Sat, 16 July 2022 081147.505 and ZK has the following metadata
Copy code
Start time                          End time                      
Fri, 15 July 2022 20:14:06.361      Sat, 16 July 2022 08:05:48.715
Sat, 16 July 2022 08:11:47.505      Sat, 16 July 2022 08:19:10.509 <- creation time is Sat, 16 July 2022 08:11:48.894
Sun, 17 July 2022 07:19:47.507      Sun, 17 July 2022 08:09:01.344
Mon, 18 July 2022 00:14:11.862      Mon, 18 July 2022 00:14:44.329
Note creation time of the segment is after the timestamp of the first record. We have seem this repeating in other tables with similar configuration
we have a few records missing that seem to match the conditions above, where "timestamp of the record" = "segment start time" < "segment creation time"
this is a pinot 0.10.0 cluster, the missing records seem to form a pattern so we wonder if this is a bug
if you need additional information to help troubleshooting this issue please let us know, we'll keep trying to figure this out
Any pointers anyone? We've been reading through the source code and it doesn't look like this could happen
m
Could it be deleted due to retention?
j
one thing I want to clarify, for the hybrid table, have you confirm the data is still there when you query just the
_REALTIME
table? I want to make sure it’s not magically being filled in by the offline side
second, have you done any sort of changes to the stream ingestion code? or is it all readily available libraries in pinot already
there is also a
pinot.server.rows_with_errors
. i’m curious if you see anything ther
s
hey Johan - it has been moved to OFFLINE in the hybrid table but we have set up another REALTIME table yesterday and it is in the REALTIME table there
correct, no changes to stream ingestion code
i can't see that metric in our prometheus i'm afraid - i can see a lot of other realtime server metrics but not that one
j
i think it’s worth getting that metric in there, or at least see where it’s published in code, and see if there’s any logs around it
my experience with pinot has been it manually maintains kafka offsets. so every message gets processed no matter what. the only way it’s missing (that i can think of) is if it failed to decode for some reason.
d
some of the previous assumptions are not quite right, we have investigated further and have more details about this issue we see in our env. From our reconciliation with the data lake we worked out the message that was missing, partition and offset were used to look at the data in the Kafka topic. We then look at which segments the row should've ended in and we found that • it seems the time end criteria was met for the mutable segment before the last record was processed • the timestamp of the record was a few fractions of second after the
realtime.segment.flush.threshold.time
- for the example below 6 hours For the example below the timestamp of the event is 1658348028492 which is Wednesday, 20 July 2022 201348.492 and the metadata of the segment where the should've be indexed is
Copy code
{
 "id": "audit_event_v001_2__14__5__20220720T1413Z",
 "simpleFields": {
   "segment.crc": "897476594",
   "segment.creation.time": "1658326428420", -> Wednesday, 20 July 2022 14:13:48.420 -> 1658348028492 - 1658326428420 = 6.00002h (slightly lafter end criteria)
   "segment.download.url": "<s3://bucket/pinot/audit_event_v001_2/audit_event_v001_2__14__5__20220720T1413Z>",
   "segment.end.time": "1658348026792",     -> Wednesday, 20 July 2022 20:13:46.792
   "segment.flush.threshold.size": "312500",
   "segment.index.version": "v3",
   "segment.realtime.endOffset": "58929",   -> the offset of the missing record is actually 58928
   "segment.realtime.numReplicas": "2",
   "segment.realtime.startOffset": "58598",
   "segment.realtime.status": "DONE",
   "segment.start.time": "1658326434861",   -> Wednesday, 20 July 2022 14:13:54.861
   "segment.time.unit": "MILLISECONDS",
   "segment.total.docs": "330"            -> 58929 - 58598 = 331 (end offset is exclusive so this number is correct)
 },
 "mapFields": {},
 "listFields": {}
}
this is the same scenario for all the missing records we have seen so far so we have setup to more tables, one with low rows threshold and high time threshold and another table with high rows threshold and low time threshold and they have worked fine so far. I'll add more details about some code paths that might be causing the issue
@Mayank apologies for the mention on this thread but I was wondering if this problem has been seen before by other users
This is what we think could be happening 1) consumeLoop() while loop https://github.com/apache/pinot/blob/release-0.10.0/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/LLRealtimeSegmentDataManager.java#L394 2) processStreamEvents() function https://github.com/apache/pinot/blob/release-0.10.0/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/LLRealtimeSegmentDataManager.java#L420 3) processStreamEvents() for loop https://github.com/apache/pinot/blob/release-0.10.0/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/LLRealtimeSegmentDataManager.java#L475 4) advance condition https://github.com/apache/pinot/blob/release-0.10.0/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/LLRealtimeSegmentDataManager.java#L422 5) discard unlfiltered messages https://github.com/apache/pinot/blob/release-0.10.0/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/LLRealtimeSegmentDataManager.java#L432 The table needs to have a segment time threshold that could be reached before the max row count threshold is reached, e.g. low traffic topic with big row count threshold Example:
LLRealtimeSegmentDataManager
enters loop 1) because end criteria is not met and processes the first batch, calls
processStreamEvents()
at line 2). Inside this function the end criteria check at 3) is still not met and records are indexed, this means
_currentOffset
has been updated. At this point condition at 4) is met and
lastUpdatedOffset
is updated with
_currentOffset
The data manager enters loop 1) again and the end criteria is not met, reads the next batch and calls
processStreamEvents()
at line 2) again. This time the end criteria check at 3) is met because time has elapsed beyond the segment time threshold and no records are processed inside this function, this means
_currentOffset
hasn't been updated. At this point condition at 4) is not met because
lastUpdatedOffset
is equal to
_currentOffset
. Instead condition at 5) is met and
_currentOffset
is updated with the last offset in the batch, skipping all the records in the unprocessed batch. Graphic example if this makes sense
Copy code
offset   current offset   lastupdated offset    batch # (start,last)
0        0                0                      1 (0,1)
1        1                0                      1 (0,1)     -> _currentOffset.compareTo(lastUpdatedOffset) != 0
2        2                2                      2 (2,2)
=== end criteria reached inside processStreamEvents()        -> messageBatch.getUnfilteredMessageCount() > 0 => lastUpdatedOffset = _currentOffset = nextOffset = 2+1
3
4
👀 1
j
hmmm, this is really interesting, and you’re making me wonder if we’re facing this as well on our clusters since we filter out a lot of events. this might be easier if you file an issue at this point
I was actually just looking at this code yesterday because we’re not implementing it in our own plugin, and I made a ticket to verify the filtered vs unfiltered logic was doing what was expected
But I noticed weird behavior where restarting a server with the “wait for consumption to catchup” feature would sometimes never catchup on certain partitions with really low volume and filtering, so I’m wondering if it’s related
d
We should raise an issue shortly. I'll try to recreate the issue in a unit test
m
@Jackie ^^
j
@Dan DC This is great analysis! I see you are adding an unit test now, do you want to also contribute the fix?
👍 1
I'm currently working on https://github.com/apache/pinot/issues/9014 which is related to this issue, so I can also fix them together, and we can use the unit test you are working on to guard both behaviors
Let me know which way do you prefer
d
Hi @Jackie it would be nice if I could contribute the fix but since you are already working on something similar it may be faster if you commit It with your change. I'll leave it up to you, I don't mind either way as long as it gets fixed :)
j
@Dan DC Sure, go ahead and contribute a fix since you already found the root cause. You may have the fix alone with the unit test in one PR
👍 1