We run Pinot 0.8.0. When ingesting a table in `FUL...
# troubleshooting
m
We run Pinot 0.8.0. When ingesting a table in
FULL
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
Copy code
| 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?
m
@Yupeng Fu ^^
I suspect the 789 might be because of partial result or some other reason not necessarily related to upsert.
y
can you check if the ingested topic is partitioned correctly?
m
the topic only has 1 partition and 1 replica per partition
789 is one example. Sometimes it could be 999 or 998, which i doubt is due to a partial result
Could it be a regression in 0.8.0 because of the partial update support?
Rolled it back to 0.7.1 in dev, the numbers returned are a bit more consistent (?), which are 1000, seldom 999 or 998, and rarely 95x. Don’t see extremes like 2000 or 789 yet.
It seems the index updates are done in batches, and the primary keys are first removed and then added, which lead to inconsistent results? This can have serious consequences if we cannot rely on the data snapshot in Pinot…
@Yupeng Fu Any insights on this?
y
if there are duplicates, that means the upsert is not configured right
select the virtual column
$hostname
and see which hosts return the duplicates
m
@Yupeng Fu we only have 1 host for this testing in dev
y
and
$segmentName
to see which segments
m
@Yupeng Fu they are from different segments
y
it'll be helpful to see your table config and schema
cc @Jackie
m
@Yupeng Fu here is the minimal tableconfig:
Copy code
{
  "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"
    }
  }
}
Also, how do you explain the case where the returned rows less than 1000? @Yupeng Fu
Actually, it seems that you will only get duplicates (>1000 rows) when querying via Trino (when a segment build is triggered?). When query Pinot directly, you only get rows less than 1000.
@Elon thoughts regarding the duplicates returned when querying via Trino?
e
Do you get correct results if you do a "passthrough" query? i.e.
Copy code
select * from pinot.default."select count(*) from <table>"
m
Message feed has stopped publishing so will have to test tomorrow. Just wanted to point it out that it’s not limited to
count()
- a regular
select * from table1
query without pushdowns can return more than and fewer than 1000 rows.
Good morning @Yupeng Fu Could you help us understand how upsert guarantees that there will always 1000 rows returned, especially during index updates?
y
m
Based on the tests so far in Trino,
select 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.
@Yupeng Fu we did already, but it’s not clear on this guarantee. Would you mind elaborating?
We suspect this has something to do with the segment builds so we did another testing with a gigantic segment to make sure a build won’t be triggered during the test. The result is that `select count(*) from table1`in PQL returns a number oscillates between 1000 and 999.
numDocsScanned reported by Pinot is also 1000 and 999 respectively
Seems the primaryKeyIndex or the validDocIndex is the culprit?
e
What is the version of trino you are using?
m
Trino version is 362
y
hmm, it shall not be, the indexing building during segment load and update is atomic
cc @Jackie
m
Telling from this pseudocode in the desgin doc,
remove
and
put/add
don’t happen within one transaction and queries may very well return different results depending on the timing?
Copy code
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)
@Yupeng Fu this can happen to a table with only a single consuming segment
j
@Map The remove/add happens in 2 steps, where we first remove (invalidate) the old record, then add the new one. Also, when executing a query, the
validDocIndex
is read when the segment is executed, so small inconsistency during record update is expected.
But the extreme values are not expected. Can you reproduce it?
m
@Jackie there is no notion or use of transactions?
Extreme values can be reproduced when you have segments being built
even one record off can cause serious problems in some use cases
j
No, pinot does not do transactions / global synchronization due to performance reasons
Pinot is mainly for analytics purpose, and is providing eventual consistency guarantee. What use cases are you referring to? Pinot should not be used as a transactional db
For the extreme value issue, could you please try with the latest release
0.9
? The upsert is a relatively new feature, and there are some changes / bugfixes recently
m
@Jackie IMHO, this is not about strong or weak consistency. Instead, it is about correctness. It’s okay to have eventual consistency, but in between, you can still return the old value for the record to be updated instead of nothing.
j
Pinot is an append only data store as of now, and upsert is handled as marking the old record invalid and add the new record. The challenge here is that there is no easy way to make these 2 updates atomic if they are applied to 2 different segments without locking all the segments, which could cause huge performance degradation. Either we add new record first or invalidate old record first, we will see duplicate key or missing key, and the current implementation is to invalidate the old record first
Here is the javadoc within the implementation describing the transient inconsistency scenarios:
Copy code
* <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>
I just submitted a PR to make inner segment update atomic, but making cross segment update atomic is very hard
👍 3
😮 1
m
Thanks for the quickaround on the PR @Jackie! Very much appreciated. I think this should solve the 99x case within one segment.
Do you have an idea on why we would get numbers like 789 or 1567 as well? I will open up a GH issue tomorrow documenting the steps to reproduce
It seems to happen between a segment finishing consuming and a new segment starting to consume
Could it be due to that old records get removed from and news ones get added into the
validDocIndex
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.
@Mayank @Jackie @Yupeng Fu @Elon I’ve opened up a GH issue https://github.com/apache/pinot/issues/7849 detailing the steps to reproduce. In summary, there are 3 issues we have discussed yesterday: • one or two rows are missing • hundreds or thousands of rows are missing • duplicates are returned The first should be addressed by https://github.com/apache/pinot/pull/7844 (Thanks @Jackie) but the other two need further investigation.
@Jackie do you think we can get https://github.com/apache/pinot/pull/7844 merged today?
m
Thanks @Map
j
Yeah, will merge it soon
❤️ 2
e
Nice!
m
@Jackie it seems you push a image build every day to dockerhub based on the latest commit. Do you know when we can expect such an image including your patch?
j
@Xiang Fu ^^
x
You can check when the last image is pushed, and expected 24 hours time period. Usually 2pm PDT.
👍 1
m
Good morning @Jackie, we tested your patch but we can still have this off-by-one row count intra-segment
Could it be that although the
vliadDocsIndex
is up-to-date, the latest row hasn’t been added?
If so, it is a regression from the partial upsert…`addNewRow()` used to happen before
handleUpsert
: https://github.com/apache/pinot/pull/6899/files#diff-a530efdea0e8088ddbe4daf9a702d57999dd43a1f140c275643831d80b4dc369
j
@Map Good finding! If you only have one segment per partition, then this should be the case
Are you running the latest master? Do you still observe large difference occasionally?
m
we are running 0.10.0-SNAPSHOT-0290d74fe-20211201-jdk11 including your commit
we have been only testing with case #1 as described in https://github.com/apache/pinot/issues/7849 with a single consuming segment, no flushes so far, and in the case, we see row count can be off-by-one from to time
would like to address case #1 first before we get to case #2 where there is a large difference
@Jackie can we make
addNewRow
happen first? 😃
j
@Map It is a little bit tricky here because for partial upsert, we need to first update the record value, then add the row. Let me see if we can break it into 2 steps
m
yeah it would make sense to do it in two steps, decoupling the index updates from row additions so that additions can happen first
or maybe you can revert to the logic prior to partial upsert to keep full upsert intact and add a special handling for partial upsert
@Jackie we tried moving
handleUpsert
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?
j
Do you have one single segment per partition?
m
yes, one single segment per partition and there is only one partition
@Jackie thanks for another PR! Our internal hacky patching is similar to yours (of course without your elegant handling of the partial upsert :)). The good news is now we see the off-by-rate has decreased from 50% to 2% during a test with 3000 repeated queries. And the bad news is, as you can see, it’s not down to zero yet…
Can you think of anything else that might be causing this?
Perhaps if you can help explain what happens from end to end when a query
select count(1) from <table>
is issued and the result is returned, we can help further validate if there are other gaps.
j
For
count(*)
query, it simply reads the
validDocIds
bitmap and calculates the cardinality
Could you describe how you performed the tests?
Is the previous fix included in the test? https://github.com/apache/pinot/pull/7844
m
yes, 7844 is included
For the tests, we have a script that continuously queries Pinot via Trino and retrieves the count…after that, we tally the results
For 
count(*)
 query, it simply reads the 
validDocIds
 bitmap and calculates the cardinality
Could you point me to the code snippet?
For 
count(*)
 query, it simply reads the 
validDocIds
 bitmap and calculates the cardinality
@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 code
j
@Map
validDocIds
is a
ThreadSafeMutableRoaringBitmap
, which always take a snapshot when read, so there should be no race condition around it
m
@Jackie any ideas on what else might go wrong?
j
No idea. I feel the inner segment inconsistency should be totally eliminated with these 2 PRs
m
Did another test over the weekend. We can say for sure that there is no inconsistency when there is no message coming in. So it has to be something with the consuming or querying while consuming.
w
It sounds like the bug is caused by partial update. If so, is it possible to disable the partial update?
m
@Weixiang SunWhat makes you come to this conclusion if you don’t mind me asking? We are using the full upsert and partial upsert is disabled. @Jackie ‘s second patch last week separates them further in the code base
w
@Map Sorry, this was my impression which is clarified by you. We saw the duplicate record problem with pinot 0.8.0. The duplicate is persisting. Do you see the similar problem?
m
We see duplicates as described in case 3 at https://github.com/apache/pinot/issues/7849
Is this your case?
w
Did you check and see if the duplicate is temporary or permanent? It seems permanent to me
m
temporary for us
w
thanks @Map
m
@Weixiang Sun did you figure out your problem? We are still experiencing the issue but also out of ideas on what could be going wrong… 😞
w
@Map we have wrong configuration. The duplicate records are gone after the configuration is fixed. It seems that we have the different problem.
e
Specifically we set
numPartitions
to 1 and
numInstancesPerPartition
to 0
@Weixiang Sun - just verified again, no dups! I think this is the config we will stick with. Thanks to @Jackie!
w
Great! Thanks @Elon! But we are not eliminating the possibilities of the first two issues @Map pointed out.
➕ 1