Stuart Coleman
07/18/2022, 1:40 PMDan DC
07/18/2022, 5:24 PMDan DC
07/18/2022, 5:26 PM"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"
}
]
}Dan DC
07/18/2022, 5:27 PM{
"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"
}
}
}Dan DC
07/18/2022, 5:28 PMDan DC
07/18/2022, 5:36 PMStart 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
Dan DC
07/18/2022, 5:41 PMDan DC
07/18/2022, 5:43 PMDan DC
07/18/2022, 5:45 PMDan DC
07/19/2022, 7:34 PMMayank
Johan Adami
07/20/2022, 5:01 AM_REALTIME table? I want to make sure it’s not magically being filled in by the offline sideJohan Adami
07/20/2022, 5:01 AMJohan Adami
07/20/2022, 5:02 AMpinot.server.rows_with_errors . i’m curious if you see anything therStuart Coleman
07/20/2022, 5:44 AMStuart Coleman
07/20/2022, 5:46 AMStuart Coleman
07/20/2022, 5:53 AMJohan Adami
07/20/2022, 3:27 PMJohan Adami
07/20/2022, 3:27 PMDan DC
07/21/2022, 11:18 AMrealtime.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
{
"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 issueDan DC
07/21/2022, 11:19 AMDan DC
07/21/2022, 12:55 PMLLRealtimeSegmentDataManager 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
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
4Johan Adami
07/21/2022, 2:11 PMJohan Adami
07/21/2022, 2:11 PMJohan Adami
07/21/2022, 2:12 PMDan DC
07/21/2022, 6:27 PMDan DC
07/22/2022, 11:38 AMMayank
Jackie
07/22/2022, 7:34 PMJackie
07/22/2022, 7:40 PMJackie
07/22/2022, 7:41 PMDan DC
07/22/2022, 7:56 PMJackie
07/22/2022, 8:54 PM