Hi, here is question on partial upsert table with ...
# troubleshooting
y
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" ] }
Result snapshot (without option skip upsert / with option)
m
@Yupeng Fu ^^
y
does full upsert behave correctly? or this happens to partial only?
cc @Qiaochu Liu
y
@Yupeng Fu we tested on partial only.
q
@yelim yu will take a look
👍 1
y
@Yupeng Fu @Qiaochu Liu please let us know how it can be solved!
q
@yelim yu 1. if you don’t specify a “comparison column”, we will reply the main time column as comparison value 2. if the event with large comparison value come first, we will ignore the later event with smaller comparison value. 3. if the event the smaller comparison value come first, we will merge the records.
Copy code
if (recordInfo._comparisonValue.compareTo(currentRecordLocation.getComparisonValue()) >= 0) {
        _reuse.clear();
        GenericRow previousRecord =
            currentRecordLocation.getSegment().getRecord(currentRecordLocation.getDocId(), _reuse);
        return _partialUpsertHandler.merge(previousRecord, record);
      } else {
        LOGGER.warn(
            "Got late event for partial-upsert: {} (current comparison value: {}, record comparison value: {}), "
                + "skipping updating the record", record, currentRecordLocation.getComparisonValue(),
            recordInfo._comparisonValue);
        return record;
      }
cc: @Jackie @Yupeng Fu
j
For partial upsert, out-of-order events are ignored, and I don't see an easy way to fix this behavior since we cannot modify the records already inserted. @yelim yu Can you please link the doc which states that out-of-order events can be handled?
y
@Jackie here is link
j
@yelim yu The linked question explains the behavior of non-upsert case. @Qiaochu Liu In the upsert documentation, does it explain the behavior on out-of-order events?
q
@yelim yu the upsert wiki is attached https://docs.pinot.apache.org/basics/data-import/upsert @Jackie i can update the wiki to add explanation for out-of-order events handling and default partial upsert mergers
👍 1
m
Thanks @Qiaochu Liu, updating the doc will be great help.
y
will partial upsert out-of-order function be not updated in the system?? or is there any plan to update this?
q
@yelim yu out-of-order records will be ignored. if the record with smaller comparison column value comes later than the record with larger comparison column values. the later event will be skipped. so they won’t be updated.