This message was deleted.
# general
s
This message was deleted.
j
Hi Ashok, It sounds like your rows could be very wide, and you are not accumulating 5m rows yet to build a segment. If you think that is the case, then try dropping
maxRowsPerSegment
down to something much smaller, like 100k, and see if that works. Also drop
maxRowsInMemory
down to 50k as well. Fyi I have seen cases where the rows are extremely wide an this setting has to be dropped very low in order to produce "reasonably sized" segments. Let us know if that works for you. If it does work then let us know what the resulting segment size (in MB) and row counts are. If not, then please provide a task log that we can look at. Thanks. John
a
Copy code
2024-02-01T12:38:45,470 ERROR [[index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh]-publish] com.google.common.util.concurrent.AggregateFuture - Input Future failed with Error
java.lang.OutOfMemoryError: Direct buffer memory
Copy code
2024-02-01T12:38:45,480 ERROR [task-runner-0-priority-0] com.google.common.util.concurrent.AggregateFuture - Input Future failed with Error
java.lang.OutOfMemoryError: Direct buffer memory
Copy code
Caused by: java.lang.OutOfMemoryError: Direct buffer memory
Copy code
2024-02-01T12:38:45,532 INFO [task-runner-0-priority-0] org.apache.druid.indexing.worker.executor.ExecutorLifecycle - Task completed with status: {
  "id" : "index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh",
  "status" : "FAILED",
  "duration" : 3629767,
  "errorMsg" : "java.util.concurrent.ExecutionException: java.lang.OutOfMemoryError: Direct buffer memory\n\tat com.go...",
  "location" : {
    "host" : null,
    "port" : -1,
    "tlsPort" : -1
  }
}
j
Did you get any `Persisted rows[nnnnnn`ies in your task log? That would indicate your ingestion reached the maxRowsInMemory limit and executed a "persist" operation. If you didn't get that far, then I think dropping maxRowsInMemory down until you see that Persist message would be a good next step. I am looking for mention of that AggregateFuture message in other cases ... haven't found anything yet ...
a
let me check John
yes am... in that particular task log am seeing persisted rows in multiple time
am seeing those persisted thing in task log
you want to see me in supervisor log
John, my understanding is processed: 611535.0 processedBytes: 2586094111.0 processedWithError: 0.0 thrownAway: 0.0 unparseable: 0.0
as per supervisor statistics
the maxDirectMemory is set to 564133888
Copy code
maxMemory	564133888
totalMemory	564133888
freeMemory	383986040
usedMemory	180147848
directMemory	134217728
while using AUTO config
is that could be issue?
see am using Kafka streaming to ingest the data
usually we do 1000 / minute --- i.e. in 1 hour it will be around 60k to 70k
but now we trying to ingest 500000 to 600000 records in 1 hour approx.,
j
500-600k records is not much ... however if the records are very wide, then it could be a lot of data. How is your JVM heap size set compared with real memory on the machine? === I would still like to understand where the task is failing ... If maxRowsInMemory was set to 150k and maxRowsPerSegment set to 5m, then you should have see the "persist" entry 4-5 times in the task log ... And if it were able to successfully build a segment, then you should see a task log entry like this:
Copy code
rg.apache.druid.segment.realtime.appenderator.StreamAppenderator - Segment[datasourcename_2024-01-25T00:00:00.000Z_2024-01-26T00:00:00.000Z_2024-01-24T23:05:05.938Z_3694] of 167,576 bytes built from 11 incremental persist(s) in 972ms; pushed to deep storage in 599ms
do you have a log statement like the above?
a
Copy code
2024-02-01T12:38:16,829 INFO [task-runner-0-priority-0] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Persisted rows[14,370] and (estimated) bytes[63,542,638]
2024-02-01T12:38:16,991 INFO [[index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Flushed in-memory data for segment[ipos_tx_combine_2023-08-01T00:00:00.000Z_2023-09-01T00:00:00.000Z_2024-01-31T14:34:35.179Z_2] spill[10] to disk in [162] ms (914 rows).
2024-02-01T12:38:17,084 INFO [[index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Flushed in-memory data for segment[ipos_tx_combine_2023-06-01T00:00:00.000Z_2023-07-01T00:00:00.000Z_2024-01-31T14:34:36.290Z_2] spill[9] to disk in [90] ms (383 rows).
am getting like this
Copy code
2024-02-01T12:38:40,090 INFO [[index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh]-appenderator-merge] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Segment[ipos_tx_combine_2023-12-01T00:00:00.000Z_2024-01-01T00:00:00.000Z_2024-01-31T14:34:33.015Z_2] of 9,354,129 bytes built from 6 incremental persist(s) in 895ms; pushed to deep storage in 458ms. Load spec is: {"type":"hdfs","path":"file:/home/ubuntu/analytics/apache-druid-28.0.0/var/druid/segments/ipos_tx_combine/20231201T000000.000Z_20240101T000000.000Z/2024-01-31T14_34_33.015Z/2_5bfaa59c-778c-4ad8-bca4-5151330b3ca1_index.zip"}
2024-02-01T12:38:42,458 INFO [[index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh]-appenderator-merge] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Segment[ipos_tx_combine_2023-02-01T00:00:00.000Z_2023-03-01T00:00:00.000Z_2024-01-31T14:34:38.479Z_2] of 16,242,617 bytes built from 10 incremental persist(s) in 1,608ms; pushed to deep storage in 759ms. Load spec is: {"type":"hdfs","path":"file:/home/ubuntu/analytics/apache-druid-28.0.0/var/druid/segments/ipos_tx_combine/20230201T000000.000Z_20230301T000000.000Z/2024-01-31T14_34_38.479Z/2_2952b459-d0db-4e5f-b219-4241410ed160_index.zip"}
2024-02-01T12:38:44,358 INFO [[index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh]-appenderator-merge] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Segment[ipos_tx_combine_2023-07-01T00:00:00.000Z_2023-08-01T00:00:00.000Z_2024-01-31T14:34:35.734Z_2] of 13,196,730 bytes built from 8 incremental persist(s) in 1,281ms; pushed to deep storage in 618ms. Load spec is: {"type":"hdfs","path":"file:/home/ubuntu/analytics/apache-druid-28.0.0/var/druid/segments/ipos_tx_combine/20230701T000000.000Z_20230801T000000.000Z/2024-01-31T14_34_35.734Z/2_85112393-d9be-4eec-9a47-4a16679ecb59_index.zip"}
2024-02-01T12:38:45,465 ERROR [[index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_jolfmfoh]-publish] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Error while publishing segments for sequenceNumber[SequenceMetadata{sequenceId=0, sequenceName='index_kafka_ipos_tx_combine_c3c68f05fbc1e1f_0', assignments=[], startOffsets={KafkaTopicPartition{partition=1, topic='null', multiTopicPartition=false}=0, KafkaTopicPartition{partition=0, topic='null', multiTopicPartition=false}=0, KafkaTopicPartition{partition=2, topic='null', multiTopicPartition=false}=0}, exclusiveStartPartitions=[], endOffsets={KafkaTopicPartition{partition=1, topic='null', multiTopicPartition=false}=194014, KafkaTopicPartition{partition=0, topic='null', multiTopicPartition=false}=199744, KafkaTopicPartition{partition=2, topic='null', multiTopicPartition=false}=217777}, sentinel=false, checkpointed=true}]
java.lang.OutOfMemoryError: Direct buffer memory
	at java.nio.Bits.reserveMemory(Bits.java:175) ~[?:?]
John, to understand I am asking see .... processedBytes: 2586094111.0 says around 2GB whereas in Druid memory status it says 0.5 GB and 0.1 GB for buffer. So will it be the reason ?
or I need to reduce the size of maxRowInMemory and MaxRowPerSegment ?
please am set segmentation granularity to MONTH , typically in our production we will have 6000000 records / month from source db to sync
j
Yes I would try reducing. I suggest try maxRowsInMemory 50k and maxRowsPerSegment 100k. We can play around with these levels once we get some segments out.
a
ok well let me try today with this properties size
I will do try once again same data set ingestion using kafka with that above properties
j
Yes ... if you have very wide rows, then even 500k rows could be too large for a segment. Let us know how it goes 🙂
a
wide row you mean a single row with more column and data size of single row more than 5 MB ?
John, wide row you mean a single row with more column and data size of single row more than 5 MB ?
I have around 150 columns in this single datasource
Hello John, I changed the tuningConfig as you mentioned i.e maxRowsInMemory 50k and maxRowsPerSegment 100k. I retried the Kafka streaming ingestion... looks like now it's working fine... No TASKS failure now. Thank You so much! Regards Well I have 1 question or doubt, actually I created a Load data (Kafka stream ingestion) indexing spec without any such tuningConfig specification. Whatever I mentioned are set by default i.e. maxRowsInMemory : 150000 and maxRowsPerSegment : 5000000 So I tried with that and faced issue, Once I Increased the data records sync within 1 hour task duration 600000 there were I received such error. Now with your input and the recommendation I changed the tuningConfig value of default to as you specified, and as of now it looks like working fine. So If I do change only the maxRowsInMemory along to 50K also would worked without changing the maxRowsPerSegment ? And the wide row you meant ? is that single record data size i.e. the number of column and the size of the data ??
Please let me know.
thank you
once again
j
Hi Ashok, Glad to see it worked! 🙌 Yes single row with lots of data in it. Druid can handle large number of columns, and also large data per row ... but you just have to adjusting sizing parameters to accommodate 🙂 The two sizing parameters work at different phases of the ingestion process. @Sergio Ferragut wrote up a nice piece on this ... here is his video you can watch and it will explain the phases and how the parameters are used. https://imply.io/developer/videos/an-in-depth-look-at-streaming-ingestion/ Let us know if you have any questions on this.
a
🙏 Ok sure thank you!
John, having changed it is working... This one where I deployed and ingested is our UAT / Stage environment. We going to soon implement this in our production Production data's are there about 5 crore (50000000) I am planning to do with this (above changed) properties in production... Hope this will work for even it 5 crore data ingestion (i.e. probably within 4-5 hour about 1 crore data ingestion ) !?
i.e. ingesting the existing data form the source
And since we reduced both the size maxRowsInMemory and maxRowsPerSegment this leading to more segments than previous ! Will it compromise the performance while doing SQL query to fetch the data ? If so, will enabling compaction will help us in this case ?
s
How big are the resulting segments now?
j
+1 on Sergio's question ... please send screenshot of Segment panel if it's easier, showing MB and row counts for the segments being created. Yes, compaction was created to (among other things) consolidate the fragmentation caused by ingestion ... in this case you will have multiple ingestion tasks and they will each create their own segments. That being said, we can try to minimize the number of segments being created up front by gradually increasing the maxRowsPerSegment, and/or minimizing the number of ingestion tasks you actually need to handle the input stream. Also figure out how many segments you queries will need to access on average ... if the access is via indexed dictionary lookups, then often the segment scan times will be very quick and it's okay to have a larger number of segments, so long as they are not too much more than ideal. I have been using 2-3x as a rule of thumb here. And there are other factors which would make you want to move in one direction or the other here. Use the metrics emitted by the Druid services to inspect this and to help you make these types of decisions.
a
88 segments now previously it was 30 !
👆
Well in my project we decided to use 1 (slowly will 2 more for future requirement) datasource in druid as flat structure (combining of multiple tables data from source as single datasource ) and 1 kafka topic to real-time ingestion to that datasource. And as of now in source we have 5-6 crore records (about 10GB) For existing data ingestion (i.e. 5 crore data) right I am doing 1000 per 10 seconds via kafka to druid For real-time data ingestion 3000 / minute is designed we using druid single-server-deployment AUTO config.
in UAT & PROD environment 8 CPU 64 GB going to use for this druid node.
Please suggest config properties to work smoothly
Hope this 50K maxRowsInMemory and 100K maxRowsPerSegment will enable smooth data ingestion!
to above mentioned existing data ingestion ?

https://files.slack.com/files-pri/T0306CNUA90-F06GUUU3UTW/image.png▾

j
Those segments are pretty small, and your row size varies greatly. Looking at just MB size you could easily go 10-20x larger ... so you must be doing a lot of computation on the columns if you were having memory issues. I suggest try it again using 200k maxRowsPerSegment ... just to reduce the segment count some. After that I'm not sure ... do you mind sharing your metricsSpec and dimensionSpec? (you can obfuscate the fields if you want)
a
Sure John, will try with 200k
yeah sure dimensionSpec
no metric spec as we need raw data...
One question, actually our plan is to do ingest the existing data (6 crore data) first to this single flat structured data-source... And then we will enable near-real-time data from application to ingest to same data source [ for every 1 minute ] So my point, whether I can do ingest using maxRowsInMemory 50K and maxRowsPerSegment 100K / 200K while doing existing bulk data one-time ingestion once it is completed and published in deep-storage Whether I can change the tuningConfig again by maxRowsInMemory to 150K and maxRowsInSegment to 5000K i.e. to default value... Is it possible to do so ? will it impact anything or error occur ? But of-course in future to same datasource again If we need to do ingest 1 time a huge size data means will be complicate !
Or I will do enable the "compaction" so that the segmentation will be aligned to segment granularity size i.e. MONTH (i.e. it will be reduced), While I use ingestion spec with 50K in maxRowsInMemory and 100K / 200K or whathever applicable as maxRowsPerSegment !! since I going to have 16 cpu & 128 gb ram in production [ so enabling compaction, 4-5 lookups and 2-4 more datasource comes in future, hope the utilization of processor and memory will be taken care even in future ? ]
Another doubt, without using Kafka streaming to ingest this size of 6 crore data and if we use CSV files (1 crore per csv file) and using sortMerge to ingest to flat datasource to these many column and rows? will it work or same Direct buffer memory issue will come? I did for upto 1 lakh data using csv and sortMerge to ingest. but we given-up this approach due to other constraints and they wanted to do with Kafka inorder to manage / control via script
from chatGPT 1. `maxRowsInMemory`: ◦ This parameter controls the maximum number of rows to buffer in memory before flushing data to disk during the indexing process. A higher value increases the size of in-memory buffers, potentially improving ingestion throughput but also increasing memory usage. 2. `maxRowsPerSegment`: ◦ This parameter defines the target number of rows per Druid segment. Segments are the units of storage and query in Druid. A smaller value for
maxRowsPerSegment
results in more frequent segment creation but smaller segments, potentially improving query performance. A larger value creates larger segments, which can reduce storage overhead but may affect query performance. as per chatGPT looks like - having more frequent creation of smaller segment in size which is good to SQL query fetch performance ?? My thought is larger size segment with lesser number of segment creation is good ! which is right ?
If more segments with small size is good means 50K inMemory and 100K perSegment is good looks like
j
chatGPT has found some of the info ... but there are multiple factors to consider, so you have to find the "sweet spot" balance for your workload. In general you want enough segments to allow for parallelism based on how many cores you are running, but not too many that the overhead of opening segments for querying and doing maintenance operations on the cluster becomes troublesome. And real-time segment querying has different performance charactistics than Historical segment querying. As a starting rule of thumb I might use the square root rule, i.e. the number of Historical processing cores to process a given query should be about the quare root of the number of segments that the query will scan. So if you have a query that scans a month of data and there are 1000 segments in that month, then having 32 Historical cores might be a good place to start. But if your cluster is smaller and has only 12 cores, then try to get that segment count down to the 150 range if you can.
a
the cores you mean CPU ?
worker capacity - i.e. slot ?
Well, John today again I tried another test datasource to populate same 1,000,000 data using maxRowInMemory 50K and maxRowsPerSegment 200K ... it keep failing So again I tried with recommendation to the new datasource from the new kafka topic (i.e. for same size 1,000,000 data to ingest) i.e. 50K in Memory and 100K perSegment that too failing !!! Why so ? Note: for same properties config & for same data size 1 was successful whereas another data-source with another topic failing ?
Is it because of already 1 datasource is running ? that also utilizing the memory reason ?
I was shared you the dimension spec John in your private msg
j
HI Ashok, I saw your dimension spec, it looks like standard numeric and string dimensions. Fyi on this ... for any dimension field that you might want to filter on you should store as a string, so it can be dictionary encoded (i.e. indexed) for fast lookups. (this may change in the future, but currently this is how it is) If you are running the entire cluster on a single machine with 8 cpus and 64GB memory, you might just be very limited on resources.
a
Ok, Will increasing the JVM maxDirectMemory will help ?
since I suppose to do 6 crore existing data as one-time ingestion !
in production