本文にスキップ
AI News HubLIVE
原典の内容 · 翻訳・分析待ち6 分で読了

翻訳待ち:Scaling and Operating a Large dbt Project on Databricks: IFCO's Data Team on Performance, Visibility, and Debugging

記事の要約

AI サービスが一時的に利用できないため、復旧後に翻訳を補完します。ソース概要:AbstractIFCO runs one of the world's largest reusable packaging pools with hundreds of millions of crates and pallets...

翻訳待ち:Scaling and Operating a Large dbt Project on Databricks: IFCO's Data Team on Performance, Visibility, and Debugging
誤りを報告

訂正窓口はまだ利用できません。記事情報をコピーして保存できます。

訂正案内
本文へ

AI サービスが一時的に利用できないため、復旧後に翻訳を補完します。

Scaling and Operating a Large dbt Project on Databricks: IFCO's Data Team on Performance, Visibility, and Debugging | Databricks Blog Tune dbt incremental models so each run writes only the rows that changed, using liquid clustering, a deliberate merge strategy, and dynamic file pruning. That cut IFCO's core job runtime by over 60% and retired the nightly full refresh. Diagnose slow models from the real executed query plan, not the compiled SQL, since that is where the true causes appear. IFCO packaged that diagnosis into a repeatable skill that runs on every model. Run the dbt project as one task per model on Databricks Jobs with the open-source databricks-dbt-factory, not as one opaque job. That gives per-model visibility, targeted reruns, and enforced testing. Abstract IFCO runs one of the world's largest reusable packaging pools with hundreds of millions of crates and pallets. With over 2,000 employees worldwide, IFCO employs over 350 people in Germany, most of whom work at its Global Headquarters in Pullach, near Munich. The business is a circular pooling service: reusable plastic containers (RPCs) move fresh produce from growers and packers to distribution centers and retailers, then return to IFCO Service Centers to be washed, sorted, and sent out again, across more than 50 countries. Every crate and pallet is tracked through its life cycle, feeding the KPIs the business runs on: cycle time, loss, breakage, wash cost, and pool size. Turning billions of raw tracking events into trustworthy KPIs is hard for three reasons: there is a lot of data, some of it arrives late in ways that are hard to predict, and when it does it forces the pipeline to correct history it has already reported. This post shows how IFCO's data platform team, working with Databricks Forward Deployed Engineering, made that pipeline faster and cheaper. The transformation logic stays in dbt. It runs on Databricks, where each incremental setting maps to a concrete Delta Lake write behavior: which columns cluster the data, how much of the target table a write has to touch, and whether rows are merged or replaced. Getting those settings right, on the right data grain, cut the core semantic layer job's daily runtime by more than 60 percent and let IFCO retire a costly nightly full refresh. IFCO and the shape of the data problem A crate is picked, filled, shipped, returned, washed, and reused many times a year, so IFCO needs to know where every asset is and what has happened to it. IFCO introduced a semantic layer, which brings many different tracking signals into one governed view of asset activity: barcode scans as crates pass on the wash line, RFID reads at dock doors, and battery-powered trackers that report GPS position, nearby Bluetooth beacons, and temperature. (Throughout this post, "semantic layer" means these governed dbt models that turn raw tracking events into business KPIs) Three properties make this hard. Scale. Billions of tracking events flow in per day. Consolidation collapses the repeated pings from each asset into far fewer activity rows, but the tables the KPIs read are still large enough that rebuilding them from scratch is expensive. Unpredictable late arrivals. Most observations land within an expected window, but some feeds lag by weeks, and a few by months, on no fixed schedule. A pipeline that assumes the data it has today is the full picture of what happened yesterday will quietly report wrong history. Historic reconciliation. Asset activity is a sequence, so a late observation does not just fill a gap. Drop a scan into the middle of an asset's timeline and it changes what the pipeline already concluded about everything after it: where the asset went next, when its cycle started, which KPI bucket it landed in. A late arrival therefore forces the pipeline to recompute the state it has already published, not simply append a new row. The goal is to emit a good estimate quickly and converge it to the truth as late data lands, without reprocessing everything each night. Incremental processing at scale A dbt incremental model is, underneath, a set of Delta read and write behaviors, and most of the win came from one principle: make each run touch as few rows as possible, and cut them as early as possible. The first and biggest lever is the read itself, scanning only the files and the changed assets a run actually needs, because every row you avoid reading is a row that never reaches the expensive per-asset sorts, shuffles, and writes downstream. Each technique below is ordinary dbt config that turns into a specific Delta behavior. Cluster on the columns you filter and join on. Liquid clustering, keyed to the grain each model is queried by (for asset activity, the asset and the event date), lets the engine skip files instead of scanning them. It is what makes the next two techniques work. Choose the incremental strategy deliberately. The strategy decides how each run writes, and the choice follows from two questions: does each row have a stable key, and are you updating rows in place or replacing a group of them at once? For keyed, dedup-heavy upserts, merge is the default. Keyed on the real grain (for asset activity, asset_id and event_date_time) it does two things a bulk delete-and-reinsert cannot: The equi-join predicate turns on dynamic file pruning: the key values in the incoming batch skip target files that cannot contain a match, so the write touches only the slice it changes. (DBT_INTERNAL_DEST and DBT_INTERNAL_SOURCE are dbt's aliases for the target table and the incoming batch in the statement it generates.) Clustering on the same keys the merge matches on keeps that pruning tight. A row-hash guard, a matched_condition that compares a surrogate hash of each row, then skips rewriting rows that did not actually change, which saves writes and keeps the downstream change feed clean. delete+insert is the alternative: it deletes a whole group of rows by key and reinserts it. That is simpler when a run re-derives a group as a unit and the rows carry no stable identity to match on, at the cost of rewriting the group even where nothing changed. At very large volumes the two are worth benchmarking rather than assuming. Bound the write to a recent window. The same predicate mechanism has a second use. Instead of an equi-join for file pruning, a time bound restricts the write to recent data, so on the higher-volume upstream models the MERGE matches against a recent slice of the destination rather than the whole table: Because the predicate keys on when a row was ingested, not when the event happened, an event that is months old is still caught as long as it landed recently. The window only has to be wide enough to cover the gap between data landing and this job processing it. Set it too narrow and late data is silently skipped: it does not error, it just never gets processed. Only recompute what changed. Models scope their work to the assets touched by new or late data, identified from an ingestion watermark, and read a window wider than they write, so late events are captured without a full refresh. Keep Delta tidy. Heavy incremental tables turn on optimized writes and auto-compaction, or hand table maintenance to Predictive Optimization, so frequent merges do not leave behind a small-file read tax. The discipline is in applying these on the right grain and then confirming, from the actual query plan, that the engine really prunes rather than silently scanning. A worked example: consolidating observations into asset activity The busiest model in the semantic layer is the one that consolidates observations from every tracking technology into a single, location-aware stream per asset. It works out when an asset actually moved using window functions partitioned by asset and ordered by event time. When an observation carries no explicit location, it falls back to Databricks SQL's built-in H3 functions, which map each latitude/longitude to a hexagonal grid cell so that "same place" becomes a cheap comparison of cell IDs and their grid distance rather than repeated geographic-distance math. It was, by a wide margin, the single biggest consumer of runtime. The first move was not to optimize but to see what was actually running, and that distinction matters. dbt compile renders a model's SELECT with its refs resolved, but for an incremental model that is not the statement Databricks executes. Behind that compiled SELECT, dbt generates and runs a larger operation: temporary views, scans of the destination table, and the final write back to the table. The only way to find where time and memory go is to read the actual executed query plan, stage by stage, from query history, not the compiled SQL. Read that way, the plan was damning. The model was scanning billions of rows, spilling hundreds of gigabytes to disk, and spending about 85 percent of its time in a single per-asset window sort and shuffle. It was, in effect, rebuilding the entire table on every run. Three things caused that: The set of "changed" assets never shrank. An upstream timestamp was regenerated on every run instead of being carried through from the source, so almost every asset looked new. When everything looks dirty, an incremental run quietly becomes a full refresh. The per-asset recompute had no bound. The window functions looked all the way back through each asset's history, so even a genuinely small changed set dragged years of history through the sort. The window stage computed columns the final query never used, including a second, forward-looking window whose entire output was discarded. Each fix follows directly from its cause: carry the real ingestion timestamp through the upstream models so the changed set reflects genuinely new data, bind the recompute to a recent window, drop the unused columns and the forward window, cluster on the grain the model is queried by And finally run the whole graph as parallel per-model tasks on serverless compute (the next section). Together these cut the core job's runtime by more than 60 percent, close to two thirds, and removed the nightly full refresh that had been required to keep the KPIs correct. From one debugging session to a repeatable skill The diagnosis above: read the real executed plan instead of the compiled SQL, check the read side and the write side, trace each symptom to a root cause - is not specific to the consolidation model. It's a sequence any engineer would run on any slow incremental model on Databricks. That sequence is what gets packaged as a skill: a playbook an AI agent executes on demand, so the diagnosis scales with the number of models instead of with the number of engineers who remember how to do it. The skill mirrors the worked example step by step. It pulls the actual statement family from query history, not dbt compile output, because for an incremental model those are different statements. It reads both sides of the run: scan-side metrics (files pruned, rows read, spill) and write-side metrics (rows written vs. rows deleted), since amplification only shows up on the write side. Then it checks for the same three failure classes found in the consolidation model: a changed-set that never shrinks (an upstream timestamp regenerated instead of carried through), an unbounded per-asset recompute (a window with no lookback limit), and wasted work (columns or window passes computed but never read downstream). Each check is grounded in a metric or a plan signal, not a hunch. The output is a report, not a silent fix: each finding is stated with its evidence (rows scanned, spill bytes, plan node), paired with a proposed change, and nothing is applied to a model until it's approved. Once approved, the same before/after metrics used to justify the fix are re-measured on the next run, so the skill closes the loop instead of assuming the fix worked. The gain is consistency, not novelty. The thre [truncated for AI cost control]

要点と分析を開く

記事インテリジェンス

エンジニア上級

要点

  • AI 生成が一時的に利用できないため、ソース内容とフォールバックメタデータを保存しました。
  • AbstractIFCO runs one of the world's largest reusable packaging pools with hundreds of millions of crates and pallets...

要点と分析は自動生成され、誤りを含む場合があります。原典をご確認ください。