Map
11/29/2021, 5:23 PMFULL upsert mode, we notice the number of rows returned for the same query varies across times, but it is supposed to remain consistent.
For example, there are 1000 unique values keyed on column A, which we use as the primary key for the pinot table table1 . A query like select count(1) from table1 can return values 1567, or 789, in addition to 1000.
In the case of 2000, you can find duplicated rows with different timestamps such as
| A | currenttime |
| - | ------------ |
| a | 1:00:00 |
| a | 1:00:01 |
| b | 1:00:00 |
| b | 1:00:03 |
...
In the case of 789, many rows are simply missing…
We suspect this is related to the process of updating the index for the upserted table. Have anyone seen this before?Mayank
Mayank
Yupeng Fu
11/29/2021, 5:56 PMMap
11/29/2021, 6:12 PMMap
11/29/2021, 6:13 PMMap
11/29/2021, 6:15 PMMap
11/29/2021, 6:16 PMMap
11/29/2021, 6:19 PMMap
11/29/2021, 9:20 PMYupeng Fu
11/29/2021, 9:39 PMYupeng Fu
11/29/2021, 9:40 PM$hostname and see which hosts return the duplicatesMap
11/29/2021, 9:48 PMYupeng Fu
11/29/2021, 9:52 PM$segmentName to see which segmentsMap
11/29/2021, 9:55 PMYupeng Fu
11/29/2021, 9:59 PMYupeng Fu
11/29/2021, 9:59 PMMap
11/29/2021, 10:09 PM{
"tableName": "table1",
"schema": {
"metricFieldSpecs": [
{
"name": "B",
"dataType": "DOUBLE"
}
],
"dimensionFieldSpecs": [
{
"name": "A",
"dataType": "STRING"
}
],
"dateTimeFieldSpecs": [
{
"name": "EPOCH",
"dataType": "INT",
"format": "1:SECONDS:EPOCH",
"granularity": "1:SECONDS"
}
],
"primaryKeyColumns": [
"A"
],
"schemaName": "schema1"
},
"realtime": {
"tableName": "table1",
"tableType": "REALTIME",
"segmentsConfig": {
"schemaName": "schema1",
"timeColumnName": "EPOCH",
"replicasPerPartition": "1",
"retentionTimeUnit": "DAYS",
"retentionTimeValue": "4",
"segmentPushType": "APPEND",
"completionConfig": {
"completionMode": "DOWNLOAD"
}
},
"tableIndexConfig": {
"invertedIndexColumns": [
"A"
],
"loadMode": "MMAP",
"nullHandlingEnabled": false,
"streamConfigs": {
"realtime.segment.flush.threshold.rows": "10000",
"realtime.segment.flush.threshold.time": "96h",
"streamType": "kafka",
"stream.kafka.consumer.type": "lowLevel",
"stream.kafka.topic.name": "topic1",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.stream.kafka.KafkaJSONMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.consumer.prop.auto.offset.reset": "largest"
}
},
"tenants": {},
"metadata": {},
"routing": {
"instanceSelectorType": "strictReplicaGroup"
},
"upsertConfig": {
"mode": "FULL"
}
}
}Map
11/29/2021, 10:15 PMMap
11/29/2021, 10:21 PMMap
11/29/2021, 10:22 PMElon
11/29/2021, 10:26 PMselect * from pinot.default."select count(*) from <table>"Map
11/29/2021, 10:33 PMcount() - a regular select * from table1 query without pushdowns can return more than and fewer than 1000 rows.Map
11/30/2021, 2:30 PMYupeng Fu
11/30/2021, 4:22 PMMap
11/30/2021, 4:23 PMselect count(*) from table1 or `select * from pinot.default."select count() from table1"`returns a num of rows less than or equal to 1000, whereas `select count() from table1 where from_unixtime(B) > current_timestamp - interval '15' minute` can give a number less than, equal to, or greater than 1000.Map
11/30/2021, 4:24 PMMap
11/30/2021, 5:51 PMMap
11/30/2021, 5:57 PMMap
11/30/2021, 5:57 PMElon
11/30/2021, 5:58 PMMap
11/30/2021, 6:00 PMYupeng Fu
11/30/2021, 6:10 PMYupeng Fu
11/30/2021, 6:10 PMMap
11/30/2021, 6:12 PMremove and put/add don’t happen within one transaction and queries may very well return different results depending on the timing?
def ingestionWithUpdate(PrimaryKeyIndex sealedSegmentIndex,
PrimaryKeyIndex mutableSegmentIndex,
ValidDocIndex validDocIndex):
while msg := kafkaConsumer.poll() do:
if msg.primaryKey in mutableSegmentIndex then:
if msg.time > mutableSegmentIndex.get(msg.primaryKey).time then:
mutableSegmentIndex.put(msg.primaryKey, IndexTuple(segmentName,docId,msg.time))
updateValidDocIndex(mutableSegmentIndex.get(msg.primaryKey))
else if msg.primaryKey in sealedSegmentIndex then:
if msg.time > sealedSegmentIndex.get(msg.primaryKey).time then:
sealedSegmentIndex.remove(msg.primaryKey)
mutableSegmentIndex.put(msg.primaryKey, IndexTuple(segmentName,docId,msg.time))
updateValidDocIndex(sealedSegmentIndex.get(msg.primaryKey))
else:
mutableSegmentIndex.put(msg.primaryKey, IndexTuple(segmentName,docId,msg.time))
if time to seal:
convert mutableSegmentIndex to add to sealedSegmentIndex
update docId of the current segment in validDocIndex
def updateValidDocIndex(IndexTuple tuple):
validDocIndex.get(tuple.segmentName).remove(tuple.docId)
validDocIndex.get(active_segmentName).add(new_docId)Map
11/30/2021, 6:15 PMJackie
11/30/2021, 7:01 PMvalidDocIndex is read when the segment is executed, so small inconsistency during record update is expected.Jackie
11/30/2021, 7:02 PMMap
11/30/2021, 7:56 PMMap
11/30/2021, 7:56 PMMap
11/30/2021, 8:01 PMJackie
11/30/2021, 8:14 PMJackie
11/30/2021, 8:16 PMJackie
11/30/2021, 8:17 PM0.9? The upsert is a relatively new feature, and there are some changes / bugfixes recentlyMap
11/30/2021, 9:51 PMJackie
11/30/2021, 10:20 PMJackie
11/30/2021, 10:20 PM* <p>There will be short term inconsistency when updating the upsert metadata, but should be consistent after the
* operation is done:
* <ul>
* <li>
* When updating a new record, it first removes the doc id from the current location, then update the new location.
* </li>
* <li>
* When adding a new segment, it removes the doc ids from the current locations before the segment being added to
* the RealtimeTableDataManager.
* </li>
* <li>
* When replacing an existing segment, after the record location being replaced with the new segment, the following
* updates applied to the new segment's valid doc ids won't be reflected to the replaced segment's valid doc ids.
* </li>
* </ul>Jackie
11/30/2021, 10:21 PMMap
12/01/2021, 2:48 AMMap
12/01/2021, 2:50 AMMap
12/01/2021, 2:52 AMMap
12/01/2021, 3:10 AMvalidDocIndex but new ones haven’t been actually added into the segment? As per https://github.com/apache/pinot/blob/5333b6da3c5bcac81d05aef68032c4c30cda8402/pino[…]inot/segment/local/indexsegment/mutable/MutableSegmentImpl.java, addNewRow happens after the handleUpsert . If there are multiple threads working on these tasks concurrently, and they are pre-empted non-deterministically, a race condition can happen.Map
12/01/2021, 7:54 PMMap
12/01/2021, 7:55 PMMayank
Jackie
12/01/2021, 8:09 PMElon
12/01/2021, 8:11 PMMap
12/01/2021, 8:24 PMJackie
12/01/2021, 8:57 PMXiang Fu
Map
12/02/2021, 12:53 PMMap
12/02/2021, 3:56 PMvliadDocsIndex is up-to-date, the latest row hasn’t been added?Map
12/02/2021, 5:30 PMhandleUpsert : https://github.com/apache/pinot/pull/6899/files#diff-a530efdea0e8088ddbe4daf9a702d57999dd43a1f140c275643831d80b4dc369Jackie
12/02/2021, 6:19 PMJackie
12/02/2021, 6:19 PMMap
12/02/2021, 6:25 PMMap
12/02/2021, 6:27 PMMap
12/02/2021, 6:28 PMMap
12/02/2021, 6:31 PMaddNewRow happen first? 😃Jackie
12/02/2021, 9:20 PMMap
12/02/2021, 9:24 PMMap
12/02/2021, 11:09 PMMap
12/03/2021, 12:30 AMhandleUpsert below addNewRow , which doesn’t seem to impact full upsert, and did a test. Unfortunately, we still see off-by-ones with full upsert…anything we are missing?Jackie
12/03/2021, 12:50 AMJackie
12/03/2021, 12:53 AMMap
12/03/2021, 1:02 AMMap
12/03/2021, 1:31 AMMap
12/03/2021, 1:33 AMMap
12/03/2021, 1:33 AMselect count(1) from <table> is issued and the result is returned, we can help further validate if there are other gaps.Jackie
12/03/2021, 2:24 AMcount(*) query, it simply reads the validDocIds bitmap and calculates the cardinalityJackie
12/03/2021, 2:24 AMJackie
12/03/2021, 2:26 AMMap
12/03/2021, 2:27 AMMap
12/03/2021, 2:29 AMMap
12/03/2021, 2:37 AMForCould you point me to the code snippet?query, it simply reads thecount(*)bitmap and calculates the cardinalityvalidDocIds
Map
12/03/2021, 5:32 PMFor@Jackie /@Yupeng Fu is there a lock on this read? can `validDocIds`be updates while the cardinality is calculated? It would be very helpful if you can point me to the codequery, it simply reads thecount(*)bitmap and calculates the cardinalityvalidDocIds
Jackie
12/03/2021, 7:02 PMvalidDocIds is a ThreadSafeMutableRoaringBitmap, which always take a snapshot when read, so there should be no race condition around itMap
12/03/2021, 10:58 PMJackie
12/03/2021, 11:13 PMMap
12/06/2021, 5:00 AMWeixiang Sun
12/06/2021, 5:49 PMMap
12/06/2021, 5:54 PMWeixiang Sun
12/06/2021, 6:07 PMMap
12/06/2021, 6:09 PMMap
12/06/2021, 6:09 PMWeixiang Sun
12/06/2021, 6:11 PMMap
12/06/2021, 6:15 PMWeixiang Sun
12/06/2021, 6:22 PMMap
12/09/2021, 1:42 AMWeixiang Sun
12/09/2021, 6:12 PMElon
12/09/2021, 6:15 PMnumPartitions to 1 and numInstancesPerPartition to 0Elon
12/09/2021, 6:18 PMWeixiang Sun
12/09/2021, 6:39 PM