Hi All general question about compaction vs kafka ...
# general
i
Hi All general question about compaction vs kafka streaming tasks: Suppose we ingesting from kafka (and only from it) and then want to compact as soon as possible to improve query performance. Our kafka tasks defined to run for 1hr with granularity of 1h and queryGranularity of 5min. Compaction has skipOffsetFromLatest for 1h as well(so it wont collide with kafka ingestion tasks). However what I've noticed is that compaction not even considering time interval of last 24hrs. My theory is it's somehow related to the "lateMessageRejectionPeriod": "PT86400S" (i.e. 24hrs). So only when we start dropping events of late arrivals the compaction can starts to work on this time interval i.e. [now - 24h; now - 25h) I.e. compaction task skips not last 1hr but last 24hrs. Probably due to constant arriving of events that fall inside last 24hrs? Is my understanding correct and if it is, what can I do to make compaction actually do the job? Our situation is that most of segments remain uncompacted(500 for 1 hour within 24hrs and 12 segments for every hour outside of last 24hrs) Thanks in advance ps: we have v25
j
Hi Igor, I don't know if Druid 25 is significantly different from the current release in the area of compaction ... but assuming it is not, then yes you are correct. Ingestion/append task (e.g. streaming) lock the time intervals with a higher priority than compaction tasks, so prior to the new concurrent append/replace feature, if an interval is locked for ingestion the Coordinator won't even consider scheduling it for compaction. Streaming ingestion will place a lock on a time interval when the first record of data comes in for that interval. So if you are receiving a scattering of "late arrival" event records, theoretically the ingestion task will end up locking a lot of time intervals. The locks are released when the ingestion task cycles (e.g. taskDuration) ... but if you are running more than one ingestion task in that Supervisor, each task has a shared append lock, so the chance of the lock being completely released by all tasks would be unlikely. If the interval isn't locked when Coordinator reviews the compaction schedule it may schedule that interval for compaction, only to have the lock revoked and compaction job abort when streaming data eventually comes in for that time interval and the ingestion task needs the lock instead. Two ways I know of dealing with this: • change your __time field to use a current date/time timestamp instead of actual event date ... then the only interval being locked will be the current or prior one hour • Check out concurrent append/replace: https://druid.apache.org/docs/latest/ingestion/concurrent-append-replace ... a new feature that allows compaction to run concurrently with streaming ingestion. The new "balanced" scheduler has not been implemented for this yet, so if you use concurrent compaction you may have to orchestrate the compaction jobs yourself, otherwise Coordinator will just reschedule the latest time interval to be compacted over and over again as streaming ingestion continues to build new segments. Let us know if that helps. Thanks. John
i
Hi @John Kowtko thank you for detailed answer. Everything makes sense . for the sake of completeness I've seen also segment level locking that might be also relevant, but I also have noticed that it's not recommended for the usage and append/replace should be used instead(as far as I understand) We are looking at append/replace feature, planning to to upgrade to 29.0.1 when it will be released. I haven't completely got what you meant by 'balanced' scheduler. I assume you talking about decision of automatic compaction regarding which time intervals will be compacted (maybe you can attach some ticket references please?). I understood that we will need to schedule compaction jobs manually for those last 24hrs and probably several times.
j
> I assume you talking about decision of automatic compaction regarding which time intervals will be compacted (maybe you can attach some ticket references please?). yes. that is correct. The current algorithm is a simple "latest first" algorithm, so the Coordinator starts with the latest interval past the offset that needs compaction and isn't locked for ingestion, marks that one for compaction, and continues backwards in time. Coordinator wakes up periodically (based on indexing period?) to re-check all intervals to find the latest interval that needs compaction. With "traditional" compaction hard locks are employed, so generally intervals won't be compacted until they are past the "late arrival" range after which no more data is ingested and no more ingestion locks are being held. So once compaction runs on those intervals that are "free from locks" it will rarely have to run on them ever again. Therefore compaction can work it's way back through history and eventually compact everything. But the latest intervals are never compacted until they reach that offset or clear the always-locked intervals. With concurrent compaction however, the lockouts don't happen, so Coordinator will pick the latest interval that needs compaction, locked or not. Which generally means the latest interval. Concurrent compaction will run on that interval, compacting all historical segments that existed when the job started, and generates a new compacted set of segments. However ingestion has created one or more segments while the compaction job was running, so by the time the compaction job finishes there will be one or more new segments added to that interval, thereby needing compaction. So Coordinator will choose it again. And again ... etc. it may never leave that latest interval. There are a few things you can do to try to spread the compaction work around, such as running several compaction jobs at once, so they will take consecutive intervals. but ultimately when those jobs finish if there was any new ingestion, those intervals will have to run again. So you can also change the offsets while the jobs are running, to point the next set of jobs to a different time range ... but that is a highly manual process. The "balanced" scheduler under development will take into account "level of need", e.g. may go after intervals that have the most number of uncompacted segments, or the highest volume of uncompacted data ... so should end up taking care of the situation. It may be a compromising act though, e.g. you may end up with all of your time intervals being 90% compacted, vs 90% of your time intervals being fully compacted and 10% not compacted at all ... so I imagine it is going to take some field testing to come up with a good algorithm here that serves everyone well (... or have multiple algorithms to deal with different use cases). There is an internal design doc for this but I don't see any PRs as yet ...
i
makes total sense to me. Thanks you for very detailed answer! Great feature ahead of us. Meanwhile scheduling manually compaction is a way to go and doesn't seem like too hard to implement (and maybe even apply some heuristic based on segments metadata tables available)
j
If you end up creating a scheduling algorithm yourself based on segment metadata info, please post your results so we can see how well it works. This would make a nice white paper šŸ™‚ Thanks. John
šŸ‘ 1
i
@John Kowtko updating that we've upgraded to 29.0.1 today and I've installed manual compaction based on concurrent locks along with automatic compactions (that once again continues to disregard all intervals within non-rejected interval, i.e. it compacts all intervals [now - retention ; now - 24h] since we have "lateMessageRejectionPeriod": "PT86400S" which is 24hrs) meanwhile I've implemented simple cron that checks what interval with [now - 23h; now - 1h] has biggest number of segments and submit number of manual compaction tasks(up to some configurable limit)
just discovered that this creates additional load on hdfs deep storage. is it possible to issue
kill
task with same context of concurrent locks?