Hi Team, i am trying to implement de-duplication ...
# general
v
Hi Team, i am trying to implement de-duplication with realtime to offline flow. things i've added as below:
Copy code
Schema:

"primaryKeyColumns": [
           "count"
]
Copy code
Realtime:

"routing": {
     "instanceSelectorType": "strictReplicaGroup"
 },

    "task": {
        "taskTypeConfigsMap": {
            "RealtimeToOfflineSegmentsTask": {
                "bufferTimePeriod": "3m",
                "bucketTimePeriod": "5m",
                "schedule": "0 */1 * * * ?",
                "mergeType": "dedup",
                "maxNumRecordsPerSegment": "10"
            }
        }
    },
but its not working! do i need to add anything else? or am i doing anything wrong? Thanks, Vishal
s
Hi Vishal, can you elaborate more on the issue that you are facing? If you have a stack trace, it would be great. Also,
primaryKeyColumn
value needs to be unique for each row.
count
sounds not a good candidate for a primary key column?
v
yeah, @Seunghyun we are using count as unique key
s
@vishal can you explain a bit more in detail how it's not working?
v
i am able to see same data in offline table, lets say i am pushing count as "10" multiple times and its pushing to offline table from realtime
s
u mean without dedup?
can you post the full table config?
v
Copy code
schemas:

{
  "schemaName": "tab5",
  "dimensionFieldSpecs": [
    {
      "name": "uuid",
      "dataType": "STRING"
    }
  ],
  "metricFieldSpecs": [
    {
      "name": "count",
      "dataType": "INT"
    }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "ts",
      "dataType": "TIMESTAMP",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    }
  ],
  "primaryKeyColumns": [
    "count"
  ]
}
Copy code
realtime:

{
    "tableName": "tab5_REALTIME",
    "tableType": "REALTIME",
    "segmentsConfig": {
        "schemaName": "tab5",
        "replication": "1",
        "timeColumnName": "ts",
        "allowNullTimeValue": false,
        "replicasPerPartition": "1",
        "retentionTimeUnit": "HOURS",
        "retentionTimeValue": "1"
    },
    "tenants": {
        "broker": "DefaultTenant",
        "server": "DefaultTenant",
        "tagOverrideConfig": {}
    },
    "tableIndexConfig": {
        "createInvertedIndexDuringSegmentGeneration": false,
        "invertedIndexColumns": [],
        "noDictionaryColumns": [],
        "streamConfigs": {
            "streamType": "kafka",
            "stream.kafka.topic.name": "pinot_test",
            "stream.kafka.broker.list": "SERVERS",
            "stream.kafka.consumer.type": "lowlevel",
            "stream.kafka.consumer.prop.auto.offset.reset": "largest",
            "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": "5",
            "realtime.segment.flush.threshold.time": "1h"
        },
        "rangeIndexColumns": [],
        "rangeIndexVersion": 2,
        "autoGeneratedInvertedIndex": false,
        "sortedColumn": [],
        "bloomFilterColumns": [],
        "loadMode": "MMAP",
        "onHeapDictionaryColumns": [],
        "varLengthDictionaryColumns": [],
        "enableDefaultStarTree": false,
        "enableDynamicStarTreeCreation": false,
        "aggregateMetrics": false,
        "nullHandlingEnabled": false
    },
    "metadata": {},
    "quota": {},
    "task": {
        "taskTypeConfigsMap": {
            "RealtimeToOfflineSegmentsTask": {
                "bufferTimePeriod": "3m",
                "bucketTimePeriod": "5m",
                "schedule": "0 */1 * * * ?",
                "mergeType": "dedup",
                "maxNumRecordsPerSegment": "10"
            }
        }
    },
    "routing": {
        "instanceSelectorType": "strictReplicaGroup"
    },
    "query": {},
    "ingestionConfig": {},
    "isDimTable": false
}
Copy code
offline:

{
  "tableName": "tab5_OFFLINE",
  "tableType": "OFFLINE",
  "segmentsConfig": {
    "schemaName": "tab5",
    "replication": "1",
    "segmentPushFrequency": "HOURLY",
    "timeColumnName": "ts",
    "allowNullTimeValue": false,
    "replicasPerPartition": "1",
    "segmentPushType": "APPEND"
  },
  "tenants": {
    "broker": "DefaultTenant",
    "server": "DefaultTenant"
  },
  "tableIndexConfig": {
    "createInvertedIndexDuringSegmentGeneration": false,
    "invertedIndexColumns": [],
    "noDictionaryColumns": [],
    "rangeIndexColumns": [],
    "rangeIndexVersion": 2,
    "autoGeneratedInvertedIndex": false,
    "sortedColumn": [],
    "bloomFilterColumns": [],
    "loadMode": "MMAP",
    "onHeapDictionaryColumns": [],
    "varLengthDictionaryColumns": [],
    "enableDefaultStarTree": false,
    "enableDynamicStarTreeCreation": false,
    "aggregateMetrics": false,
    "nullHandlingEnabled": false
  },
  "metadata": {},
  "quota": {},
  "routing": {},
  "query": {},
  "ingestionConfig": {},
  "isDimTable": false
}
s
https://docs.pinot.apache.org/basics/data-import/dedup 1. I think that you need to add
dedupConfig
to compute dedup on the realtime table.
Copy code
{ 
 ...
  "dedupConfig": { 
        "dedupEnabled": true, 
        "hashFunction": "NONE" 
   }, 
 ...
}
2. Can you try to query offline table only to check the contents of the offline table and see if there's any duplicates?
Copy code
select * from tab5_OFFLINE
Pinot keeps the time boundary and computes
max(offline table data timestamp) - 1day
. because of - 1day, some offline data will be start to be used from 1 day after.
So, the key is to configure real-time table and offline table in sync
v
1. sorry, missed it, i've to add it 2. yeah i am able to see duplicates
s
hmm interesting
v
Copy code
{ 
 ...
  "dedupConfig": { 
        "dedupEnabled": true, 
        "hashFunction": "NONE" 
   }, 
 ...
}
its still same even after enabling it
image.png
s
oh actually the timestamp values are different
v
yeah, but primary key is same
s
oh yeah nvm
i need to check the dedup logic in detail then
after applying the config, can you confirm if the realtime table is at least doing the dedup as expected?
you can query
select * from tab5_REALTIME
v
no, even in realtime table its same
s
can you try to move
count
to dimensionField to see if that makes a difference?
v
okay let me try
s
oh..
also, there's a note like the following:
Copy code
An important requirement for the Pinot dedup table is to partition the input stream by the primary key. For Kafka messages, this means the producer shall set the key in the send API. If the original stream is not partitioned, then a streaming processing job (e.g. Flink) is needed to shuffle and repartition the input stream into a partitioned one for Pinot's ingestion.
have you partition data based on the primary key? if you have multiple partitions for the kafka topic and if the row with the same primary key gets produced to multiple different partitions, we also may end up having duplicated results
if you just want to quickly try out this, you can probably create the kafka topic with 1 partition to bypass the partitioning requirement
v
we didn't try partition with key, i need to try it.
thank you @Seunghyun
s
please leave the comments if that works out. I will check it tomorrow!
v
sure @Seunghyun
@Seunghyun i've tried it, and its not working
kafka partition done, and pushed data with the key,value but not solving data duplication problem
@Seunghyun we have tried kafka partition and pushing data with the key but problem is still same
s
by the way, what’s the pinot version that you’re using?
dedup is started to be supported in
0.11.0
v
0.2.5
s
0.2.5?
I don’t think that we ever produced 0.2.5
v
ohh sorry, we are using 0.10.0
thanks for correction
s
Can you bump up the pinot to
0.11.0
or you can try to configure to use upsert first
v
okay @Seunghyun, thanks
I've tried with 0.11.0 but issue is still same
@Seunghyun
s
@vishal i will try it out based on your table config/schema today and will get back to u!
@vishal I checked
0.11.0
apache release and the master’s branch
upsert work as expected in
0.11.0
but dedup looks to be not working as expected
in master’s branch i checked that both works as expected
Copy code
{
  "tableName": "tab_REALTIME",
  "tableType": "REALTIME",
  "segmentsConfig": {
    "schemaName": "tab",
    "replication": "1",
    "timeColumnName": "ts",
    "replicasPerPartition": "1",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": "5"
  },
  "tenants": {
    "broker": "DefaultTenant",
    "server": "DefaultTenant",
    "tagOverrideConfig": {}
  },
  "tableIndexConfig": {
    "createInvertedIndexDuringSegmentGeneration": false,
    "invertedIndexColumns": [],
    "noDictionaryColumns": [],
    "streamConfigs": {
      "streamType": "kafka",
      "stream.kafka.topic.name": "snlee-dedup",
      "stream.kafka.broker.list": "localhost:29092",
      "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": "5",
      "realtime.segment.flush.threshold.time": "1h"
    },
    "rangeIndexColumns": [],
    "rangeIndexVersion": 2,
    "autoGeneratedInvertedIndex": false,
    "sortedColumn": [],
    "bloomFilterColumns": [],
    "loadMode": "MMAP",
    "onHeapDictionaryColumns": [],
    "varLengthDictionaryColumns": [],
    "enableDefaultStarTree": false,
    "enableDynamicStarTreeCreation": false,
    "aggregateMetrics": false,
    "nullHandlingEnabled": false
  },
  "metadata": {},
  "quota": {},
  "routing": {
    "instanceSelectorType": "strictReplicaGroup"
  },
  "dedupConfig": {
    "dedupEnabled": true,
    "hashFunction": "NONE"
  },
  "query": {},
  "ingestionConfig": {},
  "isDimTable": false
}
Copy code
{
  "schemaName": "tab",
  "dimensionFieldSpecs": [
    {
      "name": "key",
      "dataType": "INT"
    },
    {
      "name": "uuid",
      "dataType": "STRING"
    }
  ],
  "metricFieldSpecs": [
    {
      "name": "count",
      "dataType": "INT"
    }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "ts",
      "dataType": "TIMESTAMP",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    }
  ],
  "primaryKeyColumns": [
    "key"
  ]
}
Copy code
{
  "tableName": "tab_REALTIME",
  "tableType": "REALTIME",
  "segmentsConfig": {
    "schemaName": "tab",
    "replication": "1",
    "timeColumnName": "ts",
    "replicasPerPartition": "1",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": "5"
  },
  "tenants": {
    "broker": "DefaultTenant",
    "server": "DefaultTenant",
    "tagOverrideConfig": {}
  },
  "tableIndexConfig": {
    "createInvertedIndexDuringSegmentGeneration": false,
    "invertedIndexColumns": [],
    "noDictionaryColumns": [],
    "streamConfigs": {
      "streamType": "kafka",
      "stream.kafka.topic.name": "snlee-dedup",
      "stream.kafka.broker.list": "localhost:29092",
      "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": "5",
      "realtime.segment.flush.threshold.time": "1h"
    },
    "rangeIndexColumns": [],
    "rangeIndexVersion": 2,
    "autoGeneratedInvertedIndex": false,
    "sortedColumn": [],
    "bloomFilterColumns": [],
    "loadMode": "MMAP",
    "onHeapDictionaryColumns": [],
    "varLengthDictionaryColumns": [],
    "enableDefaultStarTree": false,
    "enableDynamicStarTreeCreation": false,
    "aggregateMetrics": false,
    "nullHandlingEnabled": false
  },
  "metadata": {},
  "quota": {},
  "routing": {
    "instanceSelectorType": "strictReplicaGroup"
  },
  "upsertConfig": {
    "mode": "FULL"
  },
  "query": {},
  "ingestionConfig": {},
  "isDimTable": false
}
the above are table config/schema for dedup/upsert that i tried
1. dedup will drop the new rows with the same primary key and keep the first event 2. upsert will override the old row with the new row when the primary keys are the same
can you try to get the
upsert
to work first?
for dedup, you need to build the pinot from the master branch code base. https://docs.pinot.apache.org/developers/developers-and-contributors/code-setup#maven
v
thank you @Seunghyun, let me try
@Seunghyun how did you push the data to kafka? through flink? or direct through terminal?
s
i deployed confluent kafka locally using docker. it comes with the UI and from there u can specify the key
v
ahh okay, thanks
@Seunghyun i've installed confluent kafka and created topic but not getting idea how to push data to topic through UI
@Seunghyun if possible, can you please share me all the config you have for kafka topic, connect config and ksqlDB config?
s
Screenshot 2022-11-25 at 10.59.15 AM.png
i didn’t change any config
i just created the topic with 2 partitions
if you go to message part, you can specify value & key
i put key as the value of
key
column from pinot
Copy code
value:
{
  "key": 0,
  "uuid": "2345",
  "count": 2,
  "ts": 1669315509897
}
key: 0
e.g.
v
Don't I need to define connecter or anything else?
👍 1
which version of confluent are you using? @Seunghyun because i am not getting that produce option
its working thanks @Seunghyun
👍 1