Hi, what is the general strategy in kafka ingestio...
# general
y
Hi, what is the general strategy in kafka ingestion? something like .. if it is creating a lot of segments in an hour, run compaction every hour. Increase maxRowsInMeory, even if segment granularity is day. I see a lot of small segments.
j
Hi Youngsol, If you want to generate fewer, larger segments, then increase: • maxRowsPerSegment • maxTotalRows • taskDuration • intermediateHandoffPeriod And also reduce the number of ingestion tasks if you can. Also, pay attention to late arrival data ... if most of your ingested data has a timestamp of today, but there are a small number of events coming in for prior time intervals, then the ingestion tasks will create separate (and small) segments for those bits of data in the prior time intervals. if you have further questions about this, please post the section of your supervisor spec that contains the tuning config, and also post a screenshot showing segment sizes and counts for the current time interval, and over the past several time intervals. Thanks. John
y
Also, pay attention to late arrival data ... if most of your ingested data has a timestamp of today, but there are a small number of events coming in for prior time intervals, then the ingestion tasks will create separate (and small) segments for those bits of data in the prior time intervals.
In my case, this is the case. How to deal with late arrival data? sometimes, month late data will come due to its domain nature.
a
Do your queries rely primarily on the column chosen as the timestamp for filtering? What is your data retention strategy like? Is it possible to ingest the data using the kafka record timestamp as your primary time column, as this would be a non-decreasing key for timestamp? (And perhaps you could do range partitioning on the timestamp column that you're currently using)
y
Data retention will be 5 years, card transaction data ,which can be canceled anytime, will be stored in Druid. To overcome data deduplication (i.e approved transaction-> canceled with same transaction id), we planned to set
__time
as approved_ts as key. Deduplication process is done by windowing query.
event_sent_at
is secondary timestamp. Something looks like below. group by. __tiime window partitioned by ( transaction_id order by event_sent_at asc)
a
Does it work if you store kafka.timestamp as
__time
, ingest
approved_ts
,
event_sent_at
as they are, and change your query to:
Copy code
group by approved_ts
window partitioned by ( transaction_id order by event_sent_at asc)
I'm not certain if the above suggestion works. However, using kafka.timestamp as __time has the benefit of having a monotonic time value which will prevent segment fragmentation due to late arriving data
y
Copy code
group by approved_ts
window partitioned by ( transaction_id order by event_sent_at asc)
but it should work but I am worried that it will scan all the segments, since we do not know the ingestion time( in case of canceled transaction)
event_sent_at
is the kafka.timestamp btw.
a
but it should work but I am worried that it will scan all the segments,
If you were to set
event_sent_at
as
__time
And ingest
approved_ts
as a separate column, would you still have to scan all the segments?
Also are you already using any secondary pa_rtitioning_ for this datasource using auto compaction etc?
y
If I set
event_sent_at
as __time, in order to user
LATEST/EARLIEST
windowing function to work properly, i have to add
__time
between current_timestamp - INTERVAL ‘5’ years and current_timestamp to group by approved_ts.
a
Ah, I see.
even if segment granularity is day. I see a lot of small segments.
Was the original segment granularity DAY or did you change HOUR to DAY recently?
y
I chaged to day. It seems working now well with compaction
👍 1