Hi, here is question on partial upsert table with out of order events
When there is a timeColumnName = “update_time” which is the same as each __time column, we wanted to test whether overwrite is working under the condition that time columns is not in order. Even though we sent kafka messages not in order , we expected the overwrite result would have been worked in the order by update_time. However, after first event (actual order in the events might be the second)’s update_time is given, if the second event's (let’s assume actual order of the second event might be the first)update_time is smaller then the first event, nothing is partially overwritten since first event’s update_time is larger. We wanted to see pinot automatically set the second event as the actual first event by comparing the first event -> and automatically partial-overwrite first event “on” second event. In the documentation, it said handling out-of-order events is possible. Is it only adapted on append mode table config? if it is not, how can we use these out-of-order events in partially overwriting table?
Table config
{
"tableName":"upsertTest1",
"tableType":"REALTIME",
"segmentsConfig":{
"timeColumnName":"update_time",
"timeType":"MILLISECONDS",
"schemaName":"upsertTest1",
"replicasPerPartition":"2"
},
"tenants":{
},
"tableIndexConfig":{
"loadMode":"MMAP",
"streamConfigs":{
"streamType":"kafka",
"realtime.segment.flush.threshold.time":"6h",
"stream.kafka.consumer.type":"lowLevel",
"stream.kafka.consumer.prop.auto.offset.reset":"smallest",
"stream.kafka.topic.name":"test.upsert.test",
"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”:~~~~
},
"nullHandlingEnabled":true
},
"fieldConfigList":[
],
"metadata":{
"customConfigs":{
}
},
"routing":{
"instanceSelectorType":"strictReplicaGroup"
},
"upsertConfig":{
"mode":"PARTIAL",
"partialUpsertStrategies":{
"a":"OVERWRITE",
"b":"OVERWRITE",
"c":"OVERWRITE",
"d":"OVERWRITE",
"pending_time":"OVERWRITE",
"issued_time":"OVERWRITE",
"matched_time":"OVERWRITE",
"unmatched_time":"OVERWRITE"
}
}
}
Schema config
{
"schemaName": "upsertTest1",
"dimensionFieldSpecs": [
{
"name": "demand_id",
"dataType": "STRING"
},
{
"name": "a",
"dataType": "STRING"
},
{
"name": "b",
"dataType": "STRING"
},
{
"name": "c",
"dataType": "STRING"
},
{
"name": "d",
"dataType": "STRING"
}
],
"dateTimeFieldSpecs": [
{
"name": "update_time",
"dataType": "TIMESTAMP",
"format": "1MILLISECONDSEPOCH",
"granularity": "1:MILLISECONDS"
},
{
"name": "create_time",
"dataType": "TIMESTAMP",
"format": "1MILLISECONDSEPOCH",
"granularity": "1:MILLISECONDS"
},
{
"name": "pending_time",
"dataType": "TIMESTAMP",
"format": "1MILLISECONDSEPOCH",
"granularity": "1:MILLISECONDS"
},
{
"name": "issued_time",
"dataType": "TIMESTAMP",
"format": "1MILLISECONDSEPOCH",
"granularity": "1:MILLISECONDS"
},
{
"name": "matched_time",
"dataType": "TIMESTAMP",
"format": "1MILLISECONDSEPOCH",
"granularity": "1:MILLISECONDS"
},
{
"name": "unmatched_time",
"dataType": "TIMESTAMP",
"format": "1MILLISECONDSEPOCH",
"granularity": "1:MILLISECONDS"
}
],
"primaryKeyColumns": [
"demand_id"
]
}