Diana Arnos
12/03/2021, 4:55 PMPinot 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:
{
"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:
{
"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:
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:
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"
}Diana Arnos
12/03/2021, 5:09 PMpinot-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) [?:?]Diana Arnos
12/06/2021, 1:35 PMDiana Arnos
12/06/2021, 2:45 PMDunith Dhanushka
Diana Arnos
12/06/2021, 3:56 PMMayank
Chinmay Soman
12/06/2021, 6:49 PMYupeng Fu
12/06/2021, 7:40 PMYupeng Fu
12/06/2021, 7:41 PMQiaochu Liu
12/06/2021, 7:42 PMJackie
12/06/2021, 11:12 PMpartialUpsertStrategies. Partial upsert only update values for the columns configuredJackie
12/06/2021, 11:14 PMOVERRIDE by default when the column is not defined in the partialUpsertStrategies for partial upsertChinmay Soman
12/06/2021, 11:46 PMYupeng Fu
12/07/2021, 12:03 AMQiaochu Liu
12/07/2021, 12:05 AMDiana Arnos
12/07/2021, 8:50 AM@Diana Arnos The reason being the columns are not added to the@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. Partial upsert only update values for the columns configuredpartialUpsertStrategies
Diana Arnos
12/07/2021, 8:56 AMJackie
12/08/2021, 12:08 AMpartialUpsertStrategies from the upsertConfig
E.g.
"upsertConfig": {
"mode": "PARTIAL",
"partialUpsertStrategies": {
"formId": "OVERWRITE",
"channelId": "OVERWRITE",
"channelPlatform": "OVERWRITE",
"companyId": "OVERWRITE",
"submitted": "OVERWRITE",
"deleted": "OVERWRITE",
"createdAt": "OVERWRITE"
"deletedAt": "OVERWRITE"
},
"hashFunction": "NONE"
},Jackie
12/08/2021, 12:10 AMDiana Arnos
12/08/2021, 9:20 AMQiaochu Liu
12/08/2021, 5:36 PM(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?Jackie
12/08/2021, 6:54 PMIGNORE as a strategy where user can choose to ignore the old valueQiaochu Liu
12/08/2021, 6:59 PMDiana Arnos
12/09/2021, 8:44 AM