This message was deleted.
# general
s
This message was deleted.
j
Hi Tim, Automated materialized view support has been discussed but afaik is not on the product roadmap as of yet. The strategy of inserting new values of the same row and using "latest()" functions to retrieve the latest data is probably the more common approach taken with Druid ... combined with Druid's compaction feature which can consolidate this data automatically over time to reduce the processing burden at query time. However I have also seen use cases where the "latest" data is periodically overwritten with updated "latest" data for the current time period, resulting in a small latency for seeing the newly arriving data, but eliminating the need for the query to process it. Here are a couple of links that touch on the subject: • https://blog.hellmar-becker.de/2023/11/25/druid-data-cookbook-upserts-in-druid-sql/ • https://imply.io/developer/articles/upserts-and-data-deduplication-with-druid/
t
Hi John, thanks for your quick reply and the helpful links. I need to look deeper into the Hellmar's technique. I've seen the imply article and familiar with the technique using
latest()
, and it's something I can use. InfluxDB offers a similar approach for handling duplicates. However, the main challenge for us seems to be updating pre-aggregates – it's a critical issue. Am I correct in assuming that there isn't a straightforward solution for updating older data on ingestion time rollups across different segments without resorting to application code? I quite like the concept of ingestion time rollups, If there were just a way to automatically track and update segments when older data is added, that would eliminate our need for materialized views.
j
Hi Tim, by "update pre-aggregates" you mean updating some the original "detail" data that was aggregated in the MV? If so, then I don't know of anything other than the standard way of maintaining the MV as a separate rollup table and then periodically re-computing time slices of the MV using the latest data from the detail table. Again something which can be done fairly easily in Druid (e.g. driven via API calls) but has to be manually orchestrated, and does have the latency based on job frequency.
t
Yes, this involves updating the original "detail" raw data. For example, consider a scenario where I'm storing raw data from Kafka at 15-minute intervals. I assume that an ingestion-time monthly rollup would efficiently pre-aggregate a monthly sum every 15 minutes. Initially, everything is fine. The rollups for January and February aggregate automatically since there are no duplicate values in the below "detail" raw data table:
Copy code
Timestamp              | Insert Time          | RId | Value | Status
----------------------------------------------------------------------
2024-01-01T00:00:00Z   | 2024-01-01T00:00:00Z | R1  | 1.0   | Estimated
2024-01-01T00:15:00Z   | 2024-01-01T00:15:00Z | R1  | 4.0   | OK
...
2024-02-01T00:00:00Z   | 2024-02-01T00:00:00Z | R1  | 1.0   | OK
2024-02-01T00:15:00Z   | 2024-02-01T00:15:00Z | R1  | 2.0   | OK
However, in February (Insert Time 2024-02-01T001500Z), we receive a duplicate value for January, needing a recalculation of the January rollup. So in this scenario I assume I need re-ingest the entire January raw data using the latest values or is there other way to periodically re-compute the time slice for monthly rollup accordingly?
Copy code
Timestamp              | Insert Time          | RId | Value  | Status
----------------------------------------------------------------------
2024-01-01T00:00:00Z   | 2024-01-01T00:00:00Z | R1  | 1.0    | Estimated
2024-01-01T00:00:00Z   | 2024-02-01T00:15:00Z | R1  | 2.2    | OK         <-- NEW DUPLICATE ROW
2024-01-01T00:15:00Z   | 2024-01-01T00:15:00Z | R1  | 4.0    | OK
...
2024-02-01T00:00:00Z   | 2024-02-01T00:00:00Z | R1  | 1.0    | OK
2024-02-01T00:15:00Z   | 2024-02-01T00:15:00Z | R1  | 2.0    | OK
j
Hi Tim, my assumption was that you would maintain two tables ... the base table and the rollup table. new data. would be insert/appended into the base table, and there would be potentially two recurring jobs: 1. Rollup job updates one month at a time in the rollup table, grabbing whatever current data is in the detail table for that month, applying "latest" and other aggregations as needed. This job reads data from the base table and overwrites to the rollup table 2. Base table cleanup job updates the base data periodically consolidating newer version of existing data and removing duplicates. This job reads from and overwrites data to the base table only. In the above scenario the Rollup table is always fully aggregated so queries do not need to perform any further rollup action. However if your only aggregation is picking the latest value due to timestamp, then that is a semi-additive aggregation (doesn't sum up across time) and would not necessarily require a rollup table ... however the queries against the base table would have to perform rollup. Thanks. John
t
Hi John, you are correct in your assumption, and I appreciate your help and patience with me. I plan to store the aggregated data in separate rollup/pre-aggregate/materialized view (MV) tables (lets call it "rollup" table now on), as you described. I have a detail table for raw data at 15-minute intervals that needs to store also the duplicates as different versions, as explained in the previous message. Additionally, I will maintain three separate rollup tables for the sum of values on daily, monthly, and yearly interval. This setup needs to support both batch and streaming ingestion. However, I'm not certain if your described manual "rollup job" actually works as it is. To make this work, I think I need to track the timestamp of every inserted raw value, determine which rollup slice/segment it belongs to, and then fetch entire slice of raw data and aggregate and update the rollup tables accordingly. This seems to be totally custom implementation. Please also note that the Timestamp is __time column and Insert Time is the secondary time data. Also, I'm unsure how to set up these manual jobs in Druid ("rollup job" and "cleanup job")? I thought these types of jobs would require custom application code, as Druid doesn't support triggers, schedulers, or custom functions/procedures. My initial plan with Druid was to create separate rollup data sources from the same detailed 15-minute raw data feed, as daily/monthly/yearly rollup data sources, along with the actual detailed raw data source with 15-minute intervals. I think this would work nicely without duplicates, but updating the rollup data sources when duplicates are inserted seems quite complicated. It seems to require custom application code since I would need to identify the rollup segment each duplicate belongs to and then manually fetch raw data and insert that whole rollup segment again. https://druid.apache.org/docs/latest/tutorials/tutorial-rollup
j
Yeah. this is not an easy operation to orchestrate. Druid as a product requires you to replace data based on
__time
interval, and replace an interval in its entirety ... so if one record changes that affects a given rollup interval, unless you want to have the queries always resolving duplicates in rollups then you would need to reingest/overwrite that interval with a complete set of updated data. Since your
__time
value represents the original date of the record being affected, you would have to scan the entire detail table to look for records that have a new
insert_time
so you could process those rollup intervals. In order to avoid a full table scan you could set up a duplicate feed to go into a "rollup staging" table, put the current date/time in the
__time
field and carry forward the original
__time
value in another field ... this would provide you with a somewhat ordered queue of time intervals that require reprocessing. As for driving the batch jobs, yes you are right Druid does not have an internal job schedule as yet, although it has been discussed internally so the need is known. For now I have heard of cron and Apache Airflow as popular job scheduling tools to use. Thanks. John