So, I'm trying to configure a table with partial u...
# troubleshooting
d
So, I'm trying to configure a table with partial upsert on
Pinot 0.9.0
that consumes from a kafka topic and I'm experiencing a weird behaviour. Once the second message gets consumed, Pinot does a full upsert instead of the partial. So every field present in the second message gets updated and all the others are set to null (I believe because they are not present on the second message and the full upsert uses the default null values) Here's the table and schema configs: Schema:
Copy code
{
  "schemaName": "responseCount",
  "dimensionFieldSpecs": [
    {
      "name": "responseId",
      "dataType": "STRING"
    },
    {
      "name": "formId",
      "dataType": "STRING"
    },
    {
      "name": "channelId",
      "dataType": "STRING"
    },
    {
      "name": "channelPlatform",
      "dataType": "STRING"
    },
    {
      "name": "companyId",
      "dataType": "STRING"
    },
    {
      "name": "submitted",
      "dataType": "BOOLEAN"
    },
    {
      "name": "deleted",
      "dataType": "BOOLEAN"
    }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "operationDate",
      "dataType": "STRING",
      "format": "1:MILLISECONDS:SIMPLE_DATE_FORMAT:yyyy-MM-dd'T'HH:mm:ss.SSSZ",
      "granularity": "1:MILLISECONDS"
    },
    {
      "name": "createdAt",
      "dataType": "STRING",
      "format": "1:MILLISECONDS:SIMPLE_DATE_FORMAT:yyyy-MM-dd'T'HH:mm:ss.SSSZ",
      "granularity": "1:MILLISECONDS"
    },
    {
      "name": "deletedAt",
      "dataType": "STRING",
      "format": "1:MILLISECONDS:SIMPLE_DATE_FORMAT:yyyy-MM-dd'T'HH:mm:ss.SSSZ",
      "granularity": "1:MILLISECONDS"
    }
  ],
  "primaryKeyColumns": [
    "responseId"
  ]
}
Table:
Copy code
{
  "REALTIME": {
    "tableName": "responseCount_REALTIME",
    "tableType": "REALTIME",
    "segmentsConfig": {
      "allowNullTimeValue": false,
      "replication": "1",
      "replicasPerPartition": "1",
      "timeColumnName": "operationDate",
      "schemaName": "responseCount"
    },
    "tenants": {
      "broker": "DefaultTenant",
      "server": "DefaultTenant"
    },
    "tableIndexConfig": {
      "rangeIndexVersion": 1,
      "autoGeneratedInvertedIndex": false,
      "createInvertedIndexDuringSegmentGeneration": false,
      "loadMode": "MMAP",
      "streamConfigs": {
        "streamType": "kafka",
        "stream.kafka.topic.name": "response-count.aggregation.source",
        "stream.kafka.broker.list": "kafka:9092",
        "stream.kafka.consumer.type": "lowlevel",
        "stream.kafka.consumer.prop.auto.offset.reset": "smallest",
        "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory",
        "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.stream.kafka.KafkaJSONMessageDecoder",
        "realtime.segment.flush.threshold.rows": "0",
        "realtime.segment.flush.threshold.time": "24h",
        "realtime.segment.flush.segment.size": "100M"
      },
      "enableDefaultStarTree": false,
      "enableDynamicStarTreeCreation": false,
      "aggregateMetrics": false,
      "nullHandlingEnabled": true
    },
    "metadata": {},
    "routing": {
      "instanceSelectorType": "strictReplicaGroup"
    },
    "upsertConfig": {
      "mode": "PARTIAL",
      "partialUpsertStrategies": {
        "deleted": "OVERWRITE",
        "deletedAt": "OVERWRITE"
      },
      "hashFunction": "NONE"
    },
    "isDimTable": false
  }
}
Here's the first message consumed:
Copy code
Key: {"responseId": "52d96a0d-92ea-4103-9ea9-536252324481"}
Value:
{
  "responseId": "52d96a0d-92ea-4103-9ea9-536252324481",
  "formId": "7bd28941-f9e4-45f1-a801-5c7d647cc6cd",
  "channelId": "60d11312-0e01-48d8-acce-4871b8d2365b",
  "channelPlatform": "app",
  "companyId": "00ca0142-5634-57e6-8d44-61427ea4b13d",
  "submitted": true,
  "deleted": "false",
  "createdAt": "2021-05-21T12:55:54.000+0000",
  "operationDate": "2021-05-21T12:55:54.000+0000"
}
Here's the second message consumed:
Copy code
Key: {"responseId": "52d96a0d-92ea-4103-9ea9-536252324481"}
Value:
{
  "responseId": "52d96a0d-92ea-4103-9ea9-536252324481",
  "deleted": "true",
  "deletedAt": "2021-10-21T12:55:54.000+0000",
  "operationDate": "2021-05-21T12:55:54.000+0000"
}
If I try the exact same configs on Pinot 0.8.0 and send the exact same message to kafka, Pinot shows me an error when consuming and does not ingest the data at all. Here's the error in the server log:
Copy code
pinot-server_1        | 2021/12/03 17:06:41.902 ERROR [LLRealtimeSegmentDataManager_responseCount__0__0__20211203T1704Z] [responseCount__0__0__20211203T1704Z] Caught exception while transforming the record: {
pinot-server_1        |   "fieldToValueMap" : {
pinot-server_1        |     "formId" : "7bd28941-f9e4-45f1-a801-5c7d647cc6cd",
pinot-server_1        |     "operationDate" : "2021-05-21T12:55:54.000+0000",
pinot-server_1        |     "createdAt" : "2021-05-21T12:55:54.000+0000",
pinot-server_1        |     "companyId" : "00ca0142-5634-57e6-8d44-61427ea4b13d",
pinot-server_1        |     "deletedAt" : "null",
pinot-server_1        |     "submitted" : 1,
pinot-server_1        |     "deleted" : 0,
pinot-server_1        |     "channelPlatform" : "app",
pinot-server_1        |     "channelId" : "60d11312-0e01-48d8-acce-4871b8d2365b",
pinot-server_1        |     "responseId" : "52d96a0d-92ea-4103-9ea9-536252324481"
pinot-server_1        |   },
pinot-server_1        |   "nullValueFields" : [ ]
pinot-server_1        | }
pinot-server_1        | java.lang.NullPointerException: null
pinot-server_1        |         at org.apache.pinot.segment.local.indexsegment.mutable.MutableSegmentImpl.handleUpsert(MutableSegmentImpl.java:512) ~[pinot-all-0.8.0-jar-with-dependencies.jar:0.8.0-c4ceff06d21fc1c1b88469a8dbae742a4b609808]
pinot-server_1        |         at org.apache.pinot.segment.local.indexsegment.mutable.MutableSegmentImpl.index(MutableSegmentImpl.java:469) ~[pinot-all-0.8.0-jar-with-dependencies.jar:0.8.0-c4ceff06d21fc1c1b88469a8dbae742a4b609808]
pinot-server_1        |         at org.apache.pinot.core.data.manager.realtime.LLRealtimeSegmentDataManager.processStreamEvents(LLRealtimeSegmentDataManager.java:516) [pinot-all-0.8.0-jar-with-dependencies.jar:0.8.0-c4ceff06d21fc1c1b88469a8dbae742a4b609808]
pinot-server_1        |         at org.apache.pinot.core.data.manager.realtime.LLRealtimeSegmentDataManager.consumeLoop(LLRealtimeSegmentDataManager.java:417) [pinot-all-0.8.0-jar-with-dependencies.jar:0.8.0-c4ceff06d21fc1c1b88469a8dbae742a4b609808]
pinot-server_1        |         at org.apache.pinot.core.data.manager.realtime.LLRealtimeSegmentDataManager$PartitionConsumer.run(LLRealtimeSegmentDataManager.java:560) [pinot-all-0.8.0-jar-with-dependencies.jar:0.8.0-c4ceff06d21fc1c1b88469a8dbae742a4b609808]
pinot-server_1        |         at java.lang.Thread.run(Thread.java:829) [?:?]
Ok, looks like everytime I try a partial upsert I always get a full upsert. I tried different table configs, nothing helped. I tried to check the Test classes of Pinot and looks like the Partial Update is never tested on the class that is getting the error (MutableSegmentImpl -> https://github.com/apache/pinot/blob/release-0.9.0/pinot-segment-local/src/test/ja[…]nt/local/indexsegment/mutable/MutableSegmentImplUpsertTest.java All test configuration are setup to use the full mode
Please, let me know if you need more info
d
Diana, I've reproduced the same issue that you're seeing and I'm trying to figure out if it's a bug or incorrect config. Let me get back to you on this.
d
Ok, thanks a lot 🙂
m
Thanks @Dunith Dhanushka
c
@Yupeng Fu ^^
y
hmm this is not expected. cc @Qiaochu Liu
we shall check the unit test, and see if we have assertions on not-updated columns not being null
q
gotcha, let me take a look
j
@Diana Arnos The reason being the columns are not added to the
partialUpsertStrategies
. Partial upsert only update values for the columns configured
@Qiaochu Liu @Yupeng Fu We might want to do
OVERRIDE
by default when the column is not defined in the
partialUpsertStrategies
for partial upsert
c
+1 to that suggestion
y
makes sense, we probably can add a config for the default strategy
q
gotcha, i can take this. we can use overide strategy for columns that user don’t specific any “partialUpsertStrategies”. is this correct?
🌟 1
d
@Diana Arnos The reason being the columns are not added to the 
partialUpsertStrategies
. Partial upsert only update values for the columns configured
@Jackie but I understand a partial upsert would understand to only update the columns it is configured to and if the column is not present on the record it shouldn't set it to null and overwrite the original values. The original value that didn't receive an update should be kept as is
My biggest issue is: I need to send first a "full" record and then a partial record that will only update on Pinot the fields that have been sent. I don't want a full upsert
j
@Diana Arnos In order to get the values filled from the previous record (partial upsert semantic), currently the columns have to be configured in the
partialUpsertStrategies
from the
upsertConfig
E.g.
Copy code
"upsertConfig": {
      "mode": "PARTIAL",
      "partialUpsertStrategies": {
        "formId": "OVERWRITE",
        "channelId": "OVERWRITE",
        "channelPlatform": "OVERWRITE",
        "companyId": "OVERWRITE",
        "submitted": "OVERWRITE",
        "deleted": "OVERWRITE",
        "createdAt": "OVERWRITE"
        "deletedAt": "OVERWRITE"
      },
      "hashFunction": "NONE"
    },
For a given column, if there is no value set before, it will be kept as null until the first record containing the column is ingested
d
@Jackie thanks, that solved my problem. This is not explicit in the docs and a couple of weeks back my config was working, so I felt kinda lost.
q
@Diana Arnos yes, as @Jackie mentioned, you can update tableConfig to meet your requirement. we have a map to store the fields for partial upsert. For field in the map
Copy code
(0) if field not in the mergeStrategies map, it will always use new value even it's null.
(1) if previous record null, return new record.
(2) if previous record not null. new record null. then use old value
(3) if previous record, not null, new record not null. then merge
in your scenario, for fields that get overwrite by null values, they are (0) instead of (2). @Jackie @Yupeng Fu do you think we can optimize the default behavior for “fields not in the strategies map” so that user don’t need to configure the “IgnoreIfNull” fields?
j
@Qiaochu Liu Yes, we can put a default strategy for fields not in the strategies map. Also, we should add
IGNORE
as a strategy where user can choose to ignore the old value
👍 1
q
gotcha, thanks. i’ll create a diff today.
d
Yeah, the config solved my problem and now things are working nicely 🎉 Thank you all for all the help and explanation 😄 This will allow my team to finish setting up staging environment for Pinot and after a couple of tests we will finally go to Prod 😄
👍 3