aditi
07/31/2023, 1:24 PMNavina
07/31/2023, 1:34 PMaditi
07/31/2023, 1:44 PMNavina
08/01/2023, 5:18 AMnumber of partitions * number of replicas Each server will create a consumer for each partition that is allocated to that server.
I suspect there is some underlying issue with your setup.
1. Have you checked the pinot server logs for possible errors?
2. Do you have ingestion throttling setup? Is the consumption making progress ?
3. Have you tried the debug rest apis for the table ?
Trying to find answers for above will give a better idea of whats going onaditi
08/01/2023, 5:20 AMDo you have ingestion throttling setup? Is the consumption making progress ?
how to check this?
I was reading this blog
https://www.linkedin.com/pulse/kafka-action-part-5-consumers-advanced-config-saahas-kulkarni/
It mentions some consumer configurations like
max.poll.records:
max.partition.fetch.bytes
<http://fetch.max.wait.ms|fetch.max.wait.ms>:
fetch.min.bytes:
Do i need to change update the above conf in kafka , currently they are having default values?aditi
08/01/2023, 5:21 AM1. Have you checked the pinot server logs for possible errors?i dont see any error in logs...
Have you tried the debug rest apis for the table ?no, i will read how to do, i am new to pinot.
aditi
08/01/2023, 5:22 AMaditi
08/01/2023, 5:35 PM9661 Consumed 100041 events from (rate:4958.416/s), currentOffset=2578284, numRowsConsumedSoFar=500159, numRowsIndexedSoFar=500159
9662 Consumed 100007 events from (rate:4442.781/s), currentOffset=2378185, numRowsConsumedSoFar=300060, numRowsIndexedSoFar=300060
9663 Allocating byte array store buffer of size 249744000 for: testprofile_1aug_512MB__47__6__20230801T1630Z:bytestring.dict
9664 Consumed 100028 events from (rate:3641.8845/s), currentOffset=2178153, numRowsConsumedSoFar=100028, numRowsIndexedSoFar=100028
9665 Allocating 1048576 bytes for: testprofile_1aug_512MB__47__6__20230801T1630Z:timestamp.dict
9666 Allocating 1048576 bytes for: testprofile_1aug_512MB__31__6__20230801T1630Z:timestamp.dict
9667 Allocating byte array store buffer of size 124872000 for: testprofile_1aug_512MB__16__6__20230801T1631Z:bytestring.dict
9668 Consumed 100129 events from (rate:4557.533/s), currentOffset=2478314, numRowsConsumedSoFar=400189, numRowsIndexedSoFar=400189
9669 Allocating 1048576 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:timestamp.dict
9670 Allocating 3002712 bytes for: testprofile_1aug_512MB__47__6__20230801T1630Z:timestamp.dict
9671 Allocating 3002712 bytes for: testprofile_1aug_512MB__47__6__20230801T1630Z:bytestring.dict
9672 Consumed 25622 events from (rate:426.96216/s), currentOffset=2587624, numRowsConsumedSoFar=509499, numRowsIndexedSoFar=509499
9673 Allocating 1501356 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:timestamp.dict
9674 Allocating 1501356 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:bytestring.dict
9675 Consumed 100085 events from (rate:4105.7144/s), currentOffset=2278238, numRowsConsumedSoFar=200113, numRowsIndexedSoFar=200113
9676 Allocating 1048576 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:timestamp.dict
9677 Consumed 45913 events from (rate:765.1019/s), currentOffset=2624197, numRowsConsumedSoFar=546072, numRowsIndexedSoFar=546072
9678 Consumed 100008 events from (rate:4443.4175/s), currentOffset=2378246, numRowsConsumedSoFar=300121, numRowsIndexedSoFar=300121
9679 Allocating 1048576 bytes for: testprofile_1aug_512MB__0__6__20230801T1628Z:timestamp.dict
9680 Allocating byte array store buffer of size 249744000 for: testprofile_1aug_512MB__16__6__20230801T1631Z:bytestring.dict
9681 Consumed 92189 events from (rate:1536.2784/s), currentOffset=2570503, numRowsConsumedSoFar=492378, numRowsIndexedSoFar=492378
9682 Allocating 1048576 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:timestamp.dict
9683 Consumed 100001 events from (rate:4227.4785/s), currentOffset=2478247, numRowsConsumedSoFar=400122, numRowsIndexedSoFar=400122
9684 Allocating 3002712 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:bytestring.dict
9685 Allocating 3002712 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:timestamp.dict
9686 Consumed 26292 events from (rate:438.1927/s), currentOffset=2613916, numRowsConsumedSoFar=535791, numRowsIndexedSoFar=535791
9687 Consumed 100062 events from (rate:4566.747/s), currentOffset=2578309, numRowsConsumedSoFar=500184, numRowsIndexedSoFar=500184
9688 Consumed 33362 events from (rate:555.89435/s), currentOffset=2657559, numRowsConsumedSoFar=579434, numRowsIndexedSoFar=579434
9689 Allocating 1048576 bytes for: testprofile_1aug_512MB__16__6__20230801T1631Z:timestamp.dict
9690 Allocating 1048576 bytes for: testprofile_1aug_512MB__47__6__20230801T1630Z:timestamp.dict
9691 Consumed 33424 events from (rate:556.9274/s), currentOffset=2603927, numRowsConsumedSoFar=525802, numRowsIndexedSoFar=525802
9692 Consumed 23634 events from (rate:393.85406/s), currentOffset=2637550, numRowsConsumedSoFar=559425, numRowsIndexedSoFar=559425
9693 Consumed 42154 events from (rate:702.2974/s), currentOffset=2620463, numRowsConsumedSoFar=542338, numRowsIndexedSoFar=542338
9694 Consumed 32714 events from (rate:545.17883/s), currentOffset=2690273, numRowsConsumedSoFar=612148, numRowsIndexedSoFar=612148
9695 Consumed 33802 events from (rate:563.2352/s), currentOffset=2637729, numRowsConsumedSoFar=559604, numRowsIndexedSoFar=559604
9696 Consumed 23072 events from (rate:384.44363/s), currentOffset=2660622, numRowsConsumedSoFar=582497, numRowsIndexedSoFar=582497
9697 Consumed 29664 events from (rate:494.35056/s), currentOffset=2650127, numRowsConsumedSoFar=572002, numRowsIndexedSoFar=572002
9698 Consumed 34592 events from (rate:576.40845/s), currentOffset=2724865, numRowsConsumedSoFar=646740, numRowsIndexedSoFar=646740
9699 Allocating 1048576 bytes for: testprofile_1aug_512MB__31__6__20230801T1630Z:timestamp.dict
pinot server logs has this rate which keeps changingNavina
08/01/2023, 5:52 PMDo i need to change update the above conf in kafka , currently they are having default values?You can consider increasing the
max.partition.fetch.bytes, although you could shoot yourself in the foot by running out of memory in your servers.
What does the resource allocation look like in your servers (cpu/memory/network)?
You have 50 partitions and 30 servers. Even with replication factor = 1, each server will easily have more than 1 consumer. That should be fine as all of them are operating with default consumer configurations. Tweaking the defaults should be done with caution.aditi
08/01/2023, 5:53 PMNavina
08/01/2023, 5:53 PMpinot server logs has this rate which keeps changingare these logs for the same pinot table? I do see a range of consumption rates from ~380/s to ~5k/s. so we know it is not stuck.
Navina
08/01/2023, 5:54 PM6 core, 178 GB main memory per machine (edited)what is the xMx?
aditi
08/01/2023, 5:55 PMaditi
08/01/2023, 5:56 PMNavina
08/01/2023, 5:57 PMit is not stuck but slow, I observed that in local disk pinot write about 100GB within 3-4 min after adding a new table..this is weird and possibly nothing to do with consumption itself. is it possible that the segment build takes a long time or you are bottlenecked by disk access?
Navina
08/01/2023, 5:57 PMaditi
08/01/2023, 6:00 PM{
"REALTIME": {
"tableName": "testprofile_1aug_512MB_REALTIME",
"tableType": "REALTIME",
"segmentsConfig": {
"timeType": "MILLISECONDS",
"schemaName": "testprofile_1aug_512MB",
"replicasPerPartition": "2",
"timeColumnName": "timestamp",
"peerSegmentDownloadScheme": "http",
"minimizeDataMovement": false
},
"tenants": {
"broker": "DefaultTenant",
"server": "DefaultTenant"
},
"tableIndexConfig": {
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.consumer.type": "simple",
"stream.kafka.topic.name": "pinot_ingest",
"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": "brokers",
"realtime.segment.flush.threshold.rows": "0",
"realtime.segment.flush.threshold.segment.size": "512M",
"realtime.segment.serverUploadToDeepStore": "true",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest"
},
"optimizeDictionary": false,
"optimizeDictionaryForMetrics": false,
"noDictionarySizeRatioThreshold": 0,
"rangeIndexVersion": 2,
"autoGeneratedInvertedIndex": false,
"createInvertedIndexDuringSegmentGeneration": false,
"loadMode": "MMAP",
"enableDefaultStarTree": false,
"enableDynamicStarTreeCreation": false,
"aggregateMetrics": false,
"nullHandlingEnabled": false
},
"metadata": {
"customConfigs": {}
},
"isDimTable": false
}
}
schema:
{
"schemaName": "testprofile_1aug_512MB",
"dimensionFieldSpecs": [
{
"name": "bytestring",
"dataType": "STRING"
}
],
"dateTimeFieldSpecs": [
{
"name": "timestamp",
"dataType": "LONG",
"format": "1:MILLISECONDS:EPOCH",
"granularity": "1:MILLISECONDS"
}
]
}
Kafka records are created using:
res = ''.join(random.choices(string.ascii_uppercase +
string.digits, k=recordSize))
data={
'bytestring': res,
'timestamp': fake.random_int(min=1172854400000, max=5872854400000),
}
m=json.dumps(data)
# p.poll(1)
p.produce(topic, m.encode('utf-8'),callback=receipt)aditi
08/01/2023, 6:01 PMrecordSize is 100000 (100KB ) so string length is 100000Navina
08/01/2023, 6:04 PM"realtime.segment.flush.threshold.segment.size": "512M",
this can lead to many small segments. but I wouldn't expect that to slow down the ingestion. so just curious how you landed on this 512M value.aditi
08/01/2023, 6:05 PMNavina
08/01/2023, 6:05 PMaditi
08/01/2023, 6:07 PMtestprofile_1aug_512MB__14__6__20230801T1630Z/v3/metadata.propertiesNavina
08/01/2023, 6:11 PMaditi
08/01/2023, 6:12 PMsegment.padding.character = \u0000
segment.name = testprofile_1aug_512MB__30__5__20230801T1626Z
segment.table.name = testprofile_1aug_512MB
segment.dimension.column.names = bytestring
segment.datetime.column.names = timestamp
segment.time.column.name = timestamp
segment.total.docs = 759375
segment.start.time = 1172861874212
segment.end.time = 5872849577709
segment.time.unit = MILLISECONDS
column.bytestring.cardinality = 759375
column.bytestring.totalDocs = 759375
column.bytestring.dataType = STRING
column.bytestring.bitsPerElement = 20
column.bytestring.lengthOfEachEntry = 512
column.bytestring.columnType = DIMENSION
column.bytestring.isSorted = false
column.bytestring.hasDictionary = true
column.bytestring.isSingleValues = true
column.bytestring.maxNumberOfMultiValues = -1
column.bytestring.totalNumberOfEntries = 759375
column.bytestring.isAutoGenerated = false
column.bytestring.minValue = 0005GPKKVVWSSZKI7BX7LLG4TW87GC8P5ZA8VL7OFEJBUQQCJ3Q1B0H53HEWHLD5ORAVYLGCNC9GKYISBXIA5RPFJ0FFW4S0FC8MN0BRU0FSFQQEHKFH8KHBZM9Y1MG63QFA0L8ZX157WTH2R8POINXYPD8MPDVCF51KW3AN5ULU9YJ54RY63VWA4QXIINQECMN9LI5TN9YEETAD7WNRKYORAFLTMZT7IRDDAQXJ7SR2CEHIDXMQLLCSBQNY1BTMPGQH7YS633CRBEWSKBV6OOXQK63FWO6K5X06X0QWSAVMFVDONVPIUAWBS0XBQQBSRFWLI49E7YL5FYYJ1EN886XWKIEJ0BKAL1KG08YA9GNE8DU94417KGIK79VR4H5DPY4D2Q4GTDW3VIG7GNRYOZPL48BYO5MBDPZ8YDXQJ6GRIVH64U2W7CY7KCQE2D9ASZA4VF91RD878HRAK74IEC1B22LG492YA5VUAP00R5XGF1JDLM0O859508P3SH1K
column.bytestring.maxValue = ZZZZGQFSCVP8QDKVBFMG2T6AJ5NH6XQCXT91YPSSYGXI87CFT3YZ1FTSOXOHE9OX99ERBXYDPMLM0QPU8L5AR94RSSAQSUPMSBX1UK0TYXWLLP2E1007Z208VHNOSXH0O3S2DQ7L1HA8MH3F6O0IFB60L0LA3LW5PV8CAEBBC97WZ4B9A6QZ85TSH7MJ820RVML1WNQYOPUZFTVGRF5IJCR2W663N4NXPZ1L53P3DTIPX10CCH7GYIK2LSWLMF367RNSBCW9QR329K4DPJE12C3JN61VGT5YHH9WH6NNGMIMQWU4WB02F2OQTIILS6X32ZEU5UD9ZY5E304YFB0LYVK94OWRCSQRIMJRWVSX4AP848F0LZ27504MLDDOSWFWNNBSC18OIO7FX2LT51KN2884ZHIJS1OXZRWLZDLDPNNCN6RSF8D1UZYEWQSBEFHE6KG4HRWV57OMP6FM741SLAB373MBYAP3QQSEEW1GL5JXW5JKXCPRCUIHA42LFQCJ
column.bytestring.defaultNullValue = null
column.timestamp.cardinality = 759375
column.timestamp.totalDocs = 759375
column.timestamp.dataType = LONG
column.timestamp.bitsPerElement = 20
column.timestamp.lengthOfEachEntry = 0
column.timestamp.columnType = DATE_TIME
column.timestamp.isSorted = false
column.timestamp.hasDictionary = true
column.timestamp.isSingleValues = true
column.timestamp.maxNumberOfMultiValues = -1
column.timestamp.totalNumberOfEntries = 759375
column.timestamp.isAutoGenerated = false
column.timestamp.datetimeFormat = 1:MILLISECONDS:EPOCH
column.timestamp.datetimeGranularity = 1:MILLISECONDS
column.timestamp.minValue = 1172861874212
column.timestamp.maxValue = 5872849577709
column.timestamp.defaultNullValue = 1690907199973
segment.realtime.startOffset = 1318750
segment.realtime.endOffset = 2078125
segment.index.version = v3aditi
08/02/2023, 3:12 AMNavina
08/03/2023, 6:03 AMKartik Khare
08/03/2023, 6:05 AMKartik Khare
08/03/2023, 6:08 AMaditi
08/03/2023, 6:35 AMaditi
08/03/2023, 6:45 AMNavina
08/03/2023, 4:18 PMNavina
08/03/2023, 4:22 PMNavina
08/03/2023, 4:23 PMaditi
08/04/2023, 7:16 AMKafka startTimeMs: 1691133337824
Error registering AppInfo mbean
javax.management.InstanceAlreadyExistsException: kafka.consumer:type=app-info,id=kafka_pinot_ingest_1000-507
at com.sun.jmx.mbeanserver.Repository.addMBean(Repository.java:436) ~[?:?]
at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerWithRepository(DefaultMBeanServerInterceptor.java:1855) ~[?:?]
at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerDynamicMBean(DefaultMBeanServerInterceptor.java:955) ~[?:?]
at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerObject(DefaultMBeanServerInterceptor.java:890) ~[?:?]
at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerMBean(DefaultMBeanServerInterceptor.java:320) ~[?:?]
at com.sun.jmx.mbeanserver.JmxMBeanServer.registerMBean(JmxMBeanServer.java:522) ~[?:?]
at org.apache.pinot.shaded.org.apache.kafka.common.utils.AppInfoParser.registerAppInfo(AppInfoParser.java:64) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.shaded.org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:814) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.shaded.org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:665) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.shaded.org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:646) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.shaded.org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:626) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.plugin.stream.kafka20.KafkaPartitionLevelConnectionHandler.<init>(KafkaPartitionLevelConnectionHandler.java:64) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.plugin.stream.kafka20.KafkaStreamMetadataProvider.<init>(KafkaStreamMetadataProvider.java:54) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory.createPartitionMetadataProvider(KafkaConsumerFactory.java:43) [pinot-kafka-2.0-0.12.1-shaded.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.core.data.manager.realtime.LLRealtimeSegmentDataManager.createPartitionMetadataProvider(LLRealtimeSegmentDataManager.java:1554) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.core.data.manager.realtime.LLRealtimeSegmentDataManager.<init>(LLRealtimeSegmentDataManager.java:1407) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.core.data.manager.realtime.RealtimeTableDataManager.addSegment(RealtimeTableDataManager.java:352) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.server.starter.helix.HelixInstanceDataManager.addRealtimeSegment(HelixInstanceDataManager.java:188) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.server.starter.helix.SegmentOnlineOfflineStateModelFactory$SegmentOnlineOfflineStateModel.onBecomeOnlineFromOffline(SegmentOnlineOfflineStateModelFactory.java:161) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.pinot.server.starter.helix.SegmentOnlineOfflineStateModelFactory$SegmentOnlineOfflineStateModel.onBecomeConsumingFromOffline(SegmentOnlineOfflineStateModelFactory.java:83) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at jdk.internal.reflect.GeneratedMethodAccessor795.invoke(Unknown Source) ~[?:?]
at jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[?:?]
at java.lang.reflect.Method.invoke(Method.java:566) ~[?:?]
at org.apache.helix.messaging.handling.HelixStateTransitionHandler.invoke(HelixStateTransitionHandler.java:350) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.helix.messaging.handling.HelixStateTransitionHandler.handleMessage(HelixStateTransitionHandler.java:278) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.helix.messaging.handling.HelixTask.call(HelixTask.java:97) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at org.apache.helix.messaging.handling.HelixTask.call(HelixTask.java:49) [pinot-all-0.12.1-jar-with-dependencies.jar:0.12.1-6e235a4ec2a16006337da04e118a435b5bb8f6d8]
at java.util.concurrent.FutureTask.run(FutureTask.java:264) [?:?]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) [?:?]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) [?:?]
at java.lang.Thread.run(Thread.java:829) [?:?]Kartik Khare
08/04/2023, 7:17 AM