vishal
11/17/2022, 7:19 AMSchema:
"primaryKeyColumns": [
"count"
]
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,
VishalSeunghyun
11/17/2022, 7:32 AMprimaryKeyColumn value needs to be unique for each row. count sounds not a good candidate for a primary key column?vishal
11/17/2022, 7:33 AMvishal
11/17/2022, 7:36 AMSeunghyun
11/17/2022, 7:44 AMvishal
11/17/2022, 7:45 AMSeunghyun
11/17/2022, 7:48 AMSeunghyun
11/17/2022, 7:49 AMvishal
11/17/2022, 7:50 AMschemas:
{
"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"
]
}
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
}
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
}Seunghyun
11/17/2022, 7:56 AMdedupConfig to compute dedup on the realtime table.
{
...
"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?
select * from tab5_OFFLINESeunghyun
11/17/2022, 7:57 AMmax(offline table data timestamp) - 1day . because of - 1day, some offline data will be start to be used from 1 day after.Seunghyun
11/17/2022, 7:57 AMvishal
11/17/2022, 7:57 AMSeunghyun
11/17/2022, 7:58 AMvishal
11/17/2022, 8:00 AM{
...
"dedupConfig": {
"dedupEnabled": true,
"hashFunction": "NONE"
},
...
}
its still same even after enabling itvishal
11/17/2022, 8:01 AMSeunghyun
11/17/2022, 8:01 AMvishal
11/17/2022, 8:02 AMSeunghyun
11/17/2022, 8:02 AMSeunghyun
11/17/2022, 8:03 AMSeunghyun
11/17/2022, 8:03 AMSeunghyun
11/17/2022, 8:03 AMselect * from tab5_REALTIMEvishal
11/17/2022, 8:04 AMSeunghyun
11/17/2022, 8:07 AMcount to dimensionField to see if that makes a difference?vishal
11/17/2022, 8:07 AMSeunghyun
11/17/2022, 8:09 AMSeunghyun
11/17/2022, 8:10 AMAn 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.Seunghyun
11/17/2022, 8:11 AMSeunghyun
11/17/2022, 8:12 AMvishal
11/17/2022, 8:14 AMvishal
11/17/2022, 8:14 AMSeunghyun
11/17/2022, 8:16 AMvishal
11/17/2022, 8:17 AMvishal
11/22/2022, 10:39 AMvishal
11/22/2022, 10:40 AMvishal
11/24/2022, 7:10 AMSeunghyun
11/24/2022, 7:17 AMSeunghyun
11/24/2022, 7:17 AM0.11.0vishal
11/24/2022, 7:19 AMSeunghyun
11/24/2022, 7:19 AMSeunghyun
11/24/2022, 7:19 AMvishal
11/24/2022, 7:20 AMvishal
11/24/2022, 7:21 AMSeunghyun
11/24/2022, 8:10 AM0.11.0Seunghyun
11/24/2022, 8:11 AMSeunghyun
11/24/2022, 8:11 AMvishal
11/24/2022, 8:27 AMvishal
11/24/2022, 10:42 AMvishal
11/24/2022, 10:48 AMSeunghyun
11/24/2022, 4:21 PMSeunghyun
11/24/2022, 6:56 PM0.11.0 apache release and the master’s branchSeunghyun
11/24/2022, 6:57 PM0.11.0 but dedup looks to be not working as expectedSeunghyun
11/24/2022, 6:57 PMSeunghyun
11/24/2022, 6:57 PM{
"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
}Seunghyun
11/24/2022, 6:57 PM{
"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"
]
}Seunghyun
11/24/2022, 6:58 PM{
"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
}Seunghyun
11/24/2022, 6:59 PMSeunghyun
11/24/2022, 6:59 PMSeunghyun
11/24/2022, 7:00 PMupsert to work first?Seunghyun
11/24/2022, 7:01 PMvishal
11/25/2022, 5:24 AMvishal
11/25/2022, 5:51 AMSeunghyun
11/25/2022, 6:13 AMvishal
11/25/2022, 6:14 AMvishal
11/25/2022, 9:33 AMvishal
11/25/2022, 12:03 PMSeunghyun
11/25/2022, 6:59 PMSeunghyun
11/25/2022, 6:59 PMSeunghyun
11/25/2022, 6:59 PMSeunghyun
11/25/2022, 6:59 PMSeunghyun
11/25/2022, 7:00 PMkey column from pinotSeunghyun
11/25/2022, 7:00 PMvalue:
{
"key": 0,
"uuid": "2345",
"count": 2,
"ts": 1669315509897
}
key: 0
e.g.vishal
11/25/2022, 7:01 PMvishal
11/25/2022, 7:08 PMvishal
11/25/2022, 7:24 PMvishal
11/25/2022, 9:20 PM