<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: Ponomarev Ilya</title>
    <description>The latest articles on DEV Community by Ponomarev Ilya (@ilyario).</description>
    <link>https://dev.to/ilyario</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F139256%2F3d587d42-12bf-4ae9-a713-18692cc0dff5.jpeg</url>
      <title>DEV Community: Ponomarev Ilya</title>
      <link>https://dev.to/ilyario</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/ilyario"/>
    <language>en</language>
    <item>
      <title>Lakehouse ETL on Kubernetes: an introduction to DataFlow Operator</title>
      <dc:creator>Ponomarev Ilya</dc:creator>
      <pubDate>Sun, 19 Jul 2026 06:45:18 +0000</pubDate>
      <link>https://dev.to/ilyario/lakehouse-etl-on-kubernetes-an-introduction-to-dataflow-operator-1p41</link>
      <guid>https://dev.to/ilyario/lakehouse-etl-on-kubernetes-an-introduction-to-dataflow-operator-1p41</guid>
      <description>&lt;p&gt;Building a lakehouse means continuously landing data in Apache Iceberg: Kafka events into bronze, OLTP increments, CDC into silver. The usual stack for that is Airflow + Spark plus a pile of glue code. For the common path “read → lightly reshape → write to Iceberg,” that stack is often overkill.&lt;/p&gt;

&lt;p&gt;&lt;a href="https://github.com/dataflow-operator/dataflow" rel="noopener noreferrer"&gt;DataFlow Operator&lt;/a&gt; is a Kubernetes operator that declares ingest as a CRD: &lt;code&gt;source → transformations → sink&lt;/code&gt;. It runs the processor, handles restarts and checkpoints, and—for batch jobs—can kick off post-load steps such as Spark on Iceberg tables.&lt;/p&gt;

&lt;p&gt;Below is a short refresher on ETL, then three practical ways to land data in a lakehouse: streaming Extract, minimal streaming ETL, and batch ELT with Spark after the load.&lt;/p&gt;

&lt;h2&gt;
  
  
  What ETL is (and ELT / streaming next to it)
&lt;/h2&gt;

&lt;p&gt;&lt;strong&gt;ETL&lt;/strong&gt; means Extract, Transform, Load:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;strong&gt;Extract&lt;/strong&gt; — pull data from a source (Kafka, CDC, polling SQL).&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Transform&lt;/strong&gt; — reshape it (filter, flatten, mask fields, rename keys).&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Load&lt;/strong&gt; — write to a sink. In a lakehouse that is usually an &lt;strong&gt;Apache Iceberg&lt;/strong&gt; table via a REST Catalog (Polaris, Lakekeeper, &lt;code&gt;iceberg-rest&lt;/code&gt;) or Nessie.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;Two nearby patterns:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;Streaming ETL&lt;/strong&gt; — the same loop, continuously: messages land in bronze/silver without waiting for end-of-day.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;ELT&lt;/strong&gt; — Load into Iceberg first (often bronze), then run heavy transforms with Spark/Trino/SQL &lt;em&gt;after&lt;/em&gt; the load.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;On Kubernetes it helps when streaming and batch share the same manifest shape—only the run mode changes. The sink in the examples below is the &lt;code&gt;iceberg&lt;/code&gt; connector (REST Catalog). For a Nessie catalog with branches, use the separate &lt;code&gt;nessie&lt;/code&gt; type; the Iceberg table model is the same.&lt;/p&gt;

&lt;h2&gt;
  
  
  Where DataFlow fits
&lt;/h2&gt;

&lt;p&gt;One pipeline runtime pattern, two kinds:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Kind&lt;/th&gt;
&lt;th&gt;Mode&lt;/th&gt;
&lt;th&gt;When to use it&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;DataFlow&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Continuous Deployment&lt;/td&gt;
&lt;td&gt;Steady stream into Iceberg (Kafka, CDC)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;DataFlowCron&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;CronJob + per-tick Job&lt;/td&gt;
&lt;td&gt;Schedule and/or Spark/SQL after a successful load (&lt;code&gt;triggers&lt;/code&gt;)&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;Pipeline model:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;source → [transformations...] → sink (iceberg)
         └─ optional errors (DLQ)
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The operator owns pod lifecycle, at-least-once delivery, checkpoints for polling sources, and Secrets integration. Heavy compute (Spark, dbt, an Airflow DAG) stays outside the processor: via &lt;code&gt;DataFlowCron&lt;/code&gt; &lt;code&gt;triggers&lt;/code&gt; or a separate downstream job on lakehouse tables.&lt;/p&gt;

&lt;p&gt;Docs: &lt;a href="https://dataflow-operator.github.io/docs/" rel="noopener noreferrer"&gt;dataflow-operator.github.io/docs&lt;/a&gt;.&lt;br&gt;&lt;br&gt;
Repo: &lt;a href="https://github.com/dataflow-operator/dataflow" rel="noopener noreferrer"&gt;github.com/dataflow-operator/dataflow&lt;/a&gt;.&lt;/p&gt;


&lt;h2&gt;
  
  
  Option 1. Streaming Extract into Iceberg
&lt;/h2&gt;

&lt;p&gt;&lt;strong&gt;Goal:&lt;/strong&gt; continuously consume Kafka and append into bronze Iceberg with no transforms.&lt;/p&gt;

&lt;p&gt;Pure Extract + Load: the consumer loop, ack, and Deployment restarts are the operator’s job. You only declare source and sink.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight yaml"&gt;&lt;code&gt;&lt;span class="na"&gt;apiVersion&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;dataflow.dataflow.io/v1&lt;/span&gt;
&lt;span class="na"&gt;kind&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;DataFlow&lt;/span&gt;
&lt;span class="na"&gt;metadata&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
  &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;kafka-to-iceberg&lt;/span&gt;
&lt;span class="na"&gt;spec&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
  &lt;span class="na"&gt;source&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;kafka&lt;/span&gt;
    &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
      &lt;span class="na"&gt;brokers&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="pi"&gt;-&lt;/span&gt; &lt;span class="s"&gt;kafka:9092&lt;/span&gt;
      &lt;span class="na"&gt;topic&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;input-topic&lt;/span&gt;
      &lt;span class="na"&gt;consumerGroup&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;dataflow-group&lt;/span&gt;
  &lt;span class="na"&gt;sink&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;iceberg&lt;/span&gt;
    &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
      &lt;span class="na"&gt;catalogURI&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;https://iceberg-catalog.example.com"&lt;/span&gt;
      &lt;span class="na"&gt;warehouse&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;main&lt;/span&gt;
      &lt;span class="na"&gt;namespace&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;bronze&lt;/span&gt;
      &lt;span class="na"&gt;table&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;events&lt;/span&gt;
      &lt;span class="na"&gt;batchSize&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="m"&gt;100&lt;/span&gt;
      &lt;span class="na"&gt;autoCreateTable&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="kc"&gt;true&lt;/span&gt;
      &lt;span class="na"&gt;authenticationType&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;BEARER&lt;/span&gt;
      &lt;span class="na"&gt;tokenSecretRef&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;iceberg-catalog&lt;/span&gt;
        &lt;span class="na"&gt;key&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;token&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Use this when:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;you need a firehose from a topic into an Iceberg table;&lt;/li&gt;
&lt;li&gt;schema cleanup can wait for Spark/Trino;&lt;/li&gt;
&lt;li&gt;ingest-time transforms are not required yet.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;The same pattern works for CDC (&lt;code&gt;postgresql-cdc&lt;/code&gt;, Debezium on Kafka)—Extract stays streaming; only the source type changes; the sink remains Iceberg.&lt;/p&gt;




&lt;h2&gt;
  
  
  Option 2. Minimal streaming ETL with transformers
&lt;/h2&gt;

&lt;p&gt;&lt;strong&gt;Goal:&lt;/strong&gt; lightly reshape JSON in flight—flatten an array, add a timestamp, drop noise, mask PII—then append to Iceberg.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;DataFlow&lt;/code&gt; supports an in-process transform chain (JSONPath via gjson): &lt;code&gt;flatten&lt;/code&gt;, &lt;code&gt;timestamp&lt;/code&gt;, &lt;code&gt;filter&lt;/code&gt;, &lt;code&gt;select&lt;/code&gt;, &lt;code&gt;remove&lt;/code&gt;, &lt;code&gt;mask&lt;/code&gt;, &lt;code&gt;snakeCase&lt;/code&gt; / &lt;code&gt;camelCase&lt;/code&gt;, &lt;code&gt;debeziumUnwrap&lt;/code&gt;, &lt;code&gt;router&lt;/code&gt;, and more.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight yaml"&gt;&lt;code&gt;&lt;span class="na"&gt;apiVersion&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;dataflow.dataflow.io/v1&lt;/span&gt;
&lt;span class="na"&gt;kind&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;DataFlow&lt;/span&gt;
&lt;span class="na"&gt;metadata&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
  &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;stock-to-iceberg&lt;/span&gt;
&lt;span class="na"&gt;spec&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
  &lt;span class="na"&gt;source&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;kafka&lt;/span&gt;
    &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
      &lt;span class="na"&gt;brokers&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="pi"&gt;-&lt;/span&gt; &lt;span class="s"&gt;kafka:9092&lt;/span&gt;
      &lt;span class="na"&gt;topic&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;stock-topic&lt;/span&gt;
      &lt;span class="na"&gt;consumerGroup&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;dataflow-group&lt;/span&gt;
  &lt;span class="na"&gt;transformations&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="pi"&gt;-&lt;/span&gt; &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;flatten&lt;/span&gt;
      &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="na"&gt;field&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;rowsStock&lt;/span&gt;
    &lt;span class="pi"&gt;-&lt;/span&gt; &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;timestamp&lt;/span&gt;
      &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="na"&gt;fieldName&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;created_at&lt;/span&gt;
  &lt;span class="na"&gt;sink&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;iceberg&lt;/span&gt;
    &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
      &lt;span class="na"&gt;catalogURI&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;https://iceberg-catalog.example.com"&lt;/span&gt;
      &lt;span class="na"&gt;warehouse&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;main&lt;/span&gt;
      &lt;span class="na"&gt;namespace&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;silver&lt;/span&gt;
      &lt;span class="na"&gt;table&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;stock_items&lt;/span&gt;
      &lt;span class="na"&gt;batchSize&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="m"&gt;100&lt;/span&gt;
      &lt;span class="na"&gt;autoCreateTable&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="kc"&gt;true&lt;/span&gt;
      &lt;span class="na"&gt;authenticationType&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;BEARER&lt;/span&gt;
      &lt;span class="na"&gt;tokenSecretRef&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;iceberg-catalog&lt;/span&gt;
        &lt;span class="na"&gt;key&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;token&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;This is not a Spark/Flink replacement for lakehouse joins and heavy analytics. For “unwrap an envelope → keep the fields you need → mask a card number → land in Iceberg,” you often do not need a separate compute cluster at ingest time.&lt;/p&gt;

&lt;p&gt;A typical CDC path: Debezium on Kafka → &lt;code&gt;debeziumUnwrap&lt;/code&gt; → &lt;code&gt;select&lt;/code&gt; / &lt;code&gt;mask&lt;/code&gt; → append to a silver Iceberg table.&lt;/p&gt;




&lt;h2&gt;
  
  
  Option 3. Batch ELT: DataFlowCron + Spark after Iceberg load
&lt;/h2&gt;

&lt;p&gt;&lt;strong&gt;Goal:&lt;/strong&gt; on a schedule, pull an OLTP increment (polling SQL), land it in bronze Iceberg, then start Spark for heavy transforms on lakehouse tables.&lt;/p&gt;

&lt;p&gt;DataFlow covers &lt;strong&gt;EL&lt;/strong&gt; (Extract + Load into Iceberg); Spark owns &lt;strong&gt;T&lt;/strong&gt; after the Job succeeds. &lt;code&gt;DataFlowCron&lt;/code&gt; exposes ordered &lt;code&gt;triggers&lt;/code&gt; that start only when the processor Job completes successfully (&lt;code&gt;JobComplete&lt;/code&gt;).&lt;/p&gt;

&lt;p&gt;Important: for post-load triggers, use a &lt;strong&gt;polling source&lt;/strong&gt; (&lt;code&gt;postgresql&lt;/code&gt;, &lt;code&gt;clickhouse&lt;/code&gt;, &lt;code&gt;trino&lt;/code&gt;, &lt;code&gt;nessie&lt;/code&gt; / &lt;code&gt;iceberg&lt;/code&gt;). Kafka under cron is streaming—the Job often never reaches “source exhausted,” so triggers never run.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight yaml"&gt;&lt;code&gt;&lt;span class="na"&gt;apiVersion&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;dataflow.dataflow.io/v1&lt;/span&gt;
&lt;span class="na"&gt;kind&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;DataFlowCron&lt;/span&gt;
&lt;span class="na"&gt;metadata&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
  &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;orders-hourly-elt&lt;/span&gt;
&lt;span class="na"&gt;spec&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
  &lt;span class="na"&gt;schedule&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;0&lt;/span&gt;&lt;span class="nv"&gt; &lt;/span&gt;&lt;span class="s"&gt;*&lt;/span&gt;&lt;span class="nv"&gt; &lt;/span&gt;&lt;span class="s"&gt;*&lt;/span&gt;&lt;span class="nv"&gt; &lt;/span&gt;&lt;span class="s"&gt;*&lt;/span&gt;&lt;span class="nv"&gt; &lt;/span&gt;&lt;span class="s"&gt;*"&lt;/span&gt;
  &lt;span class="na"&gt;concurrencyPolicy&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;Forbid&lt;/span&gt;
  &lt;span class="na"&gt;checkpointSyncOnAck&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="kc"&gt;true&lt;/span&gt;
  &lt;span class="na"&gt;source&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;postgresql&lt;/span&gt;
    &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
      &lt;span class="na"&gt;connectionStringSecretRef&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;source-db&lt;/span&gt;
        &lt;span class="na"&gt;key&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;url&lt;/span&gt;
      &lt;span class="na"&gt;table&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;orders&lt;/span&gt;
      &lt;span class="na"&gt;changeTrackingColumn&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;updated_at&lt;/span&gt;
      &lt;span class="na"&gt;orderByColumn&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;id&lt;/span&gt;
      &lt;span class="na"&gt;pollInterval&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="m"&gt;30&lt;/span&gt;
      &lt;span class="na"&gt;readBatchSize&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="m"&gt;1000&lt;/span&gt;
  &lt;span class="na"&gt;sink&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="na"&gt;type&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;iceberg&lt;/span&gt;
    &lt;span class="na"&gt;config&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
      &lt;span class="na"&gt;catalogURI&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;https://iceberg-catalog.example.com"&lt;/span&gt;
      &lt;span class="na"&gt;warehouse&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;main&lt;/span&gt;
      &lt;span class="na"&gt;namespace&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;bronze&lt;/span&gt;
      &lt;span class="na"&gt;table&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;orders&lt;/span&gt;
      &lt;span class="na"&gt;batchSize&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="m"&gt;100&lt;/span&gt;
      &lt;span class="na"&gt;autoCreateTable&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="kc"&gt;true&lt;/span&gt;
      &lt;span class="na"&gt;authenticationType&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;BEARER&lt;/span&gt;
      &lt;span class="na"&gt;tokenSecretRef&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
        &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;iceberg-catalog&lt;/span&gt;
        &lt;span class="na"&gt;key&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;token&lt;/span&gt;
  &lt;span class="na"&gt;triggers&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt;
    &lt;span class="pi"&gt;-&lt;/span&gt; &lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;start-spark&lt;/span&gt;
      &lt;span class="na"&gt;image&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="s"&gt;bitnami/kubectl:latest&lt;/span&gt;
      &lt;span class="na"&gt;command&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="pi"&gt;[&lt;/span&gt;&lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;kubectl"&lt;/span&gt;&lt;span class="pi"&gt;]&lt;/span&gt;
      &lt;span class="na"&gt;args&lt;/span&gt;&lt;span class="pi"&gt;:&lt;/span&gt; &lt;span class="pi"&gt;[&lt;/span&gt;&lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;apply"&lt;/span&gt;&lt;span class="pi"&gt;,&lt;/span&gt; &lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;-f"&lt;/span&gt;&lt;span class="pi"&gt;,&lt;/span&gt; &lt;span class="s2"&gt;"&lt;/span&gt;&lt;span class="s"&gt;/manifests/spark-application.yaml"&lt;/span&gt;&lt;span class="pi"&gt;]&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;How to read this:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;Once an hour, the CronJob starts a processor.&lt;/li&gt;
&lt;li&gt;The processor reads the PostgreSQL increment from the checkpoint (&lt;code&gt;updated_at&lt;/code&gt;, &lt;code&gt;id&lt;/code&gt;) and appends into &lt;code&gt;bronze.orders&lt;/code&gt; (Iceberg).&lt;/li&gt;
&lt;li&gt;The source is exhausted → the Job succeeds → the trigger applies a &lt;code&gt;SparkApplication&lt;/code&gt; (or runs &lt;code&gt;spark-submit&lt;/code&gt; from its image)—for example bronze → silver/gold.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;Spark is not a first-class connector here; it is a generic post-step: any container with &lt;code&gt;image&lt;/code&gt; / &lt;code&gt;command&lt;/code&gt; / &lt;code&gt;args&lt;/code&gt;. Bake the SparkApplication manifest into the trigger image (the triggers API has no &lt;code&gt;volumeMounts&lt;/code&gt;). The same hook can start an Airflow DAG or refresh Trino/BI.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;checkpointSyncOnAck&lt;/code&gt; on cron helps survive at-least-once retries on Extract; merge/idempotency in the lakehouse is usually handled in the Spark job against Iceberg (MERGE / partition overwrite).&lt;/p&gt;




&lt;h2&gt;
  
  
  Which option to pick
&lt;/h2&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Need&lt;/th&gt;
&lt;th&gt;Kind&lt;/th&gt;
&lt;th&gt;Pattern&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;Continuous land into Iceberg, no reshaping&lt;/td&gt;
&lt;td&gt;&lt;code&gt;DataFlow&lt;/code&gt;&lt;/td&gt;
&lt;td&gt;Extract (option 1)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Light JSON reshape in flight → Iceberg&lt;/td&gt;
&lt;td&gt;
&lt;code&gt;DataFlow&lt;/code&gt; + &lt;code&gt;transformations&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;Streaming ETL (option 2)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Schedule + Spark after Iceberg load&lt;/td&gt;
&lt;td&gt;
&lt;code&gt;DataFlowCron&lt;/code&gt; + &lt;code&gt;triggers&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;ELT (option 3)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Steady stream &lt;em&gt;and&lt;/em&gt; post-steps&lt;/td&gt;
&lt;td&gt;two resources or an intermediate topic&lt;/td&gt;
&lt;td&gt;DataFlow → Spark/orchestrator separately&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;Product boundaries to keep in mind:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;transforms are in-process, per message/row—not distributed lakehouse joins;&lt;/li&gt;
&lt;li&gt;triggers are an ordered Job chain, not a full DAG orchestrator;&lt;/li&gt;
&lt;li&gt;delivery is at-least-once; strong idempotency on Iceberg is usually achieved in Spark (MERGE / overwrite partition).&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Next steps
&lt;/h2&gt;

&lt;ol&gt;
&lt;li&gt;Repository and issues: &lt;a href="https://github.com/dataflow-operator/dataflow" rel="noopener noreferrer"&gt;github.com/dataflow-operator/dataflow&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;Documentation: &lt;a href="https://dataflow-operator.github.io/docs/" rel="noopener noreferrer"&gt;dataflow-operator.github.io/docs&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;Helm quick start: &lt;a href="https://dataflow-operator.github.io/docs/getting-started/" rel="noopener noreferrer"&gt;Getting started&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;Iceberg connector: &lt;a href="https://dataflow-operator.github.io/docs/connectors/" rel="noopener noreferrer"&gt;Connectors&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;Transformations: &lt;a href="https://dataflow-operator.github.io/docs/transformations/" rel="noopener noreferrer"&gt;Transformations&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;Cron and post-load: &lt;a href="https://dataflow-operator.github.io/docs/dataflow-cron/triggers/" rel="noopener noreferrer"&gt;DataFlowCron triggers&lt;/a&gt;
&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;If you already run “simple” ingest as scripts—or a heavy Airflow DAG just to move Kafka into Iceberg—consider declaring it as a &lt;code&gt;DataFlow&lt;/code&gt; or &lt;code&gt;DataFlowCron&lt;/code&gt; and keep the orchestrator for the parts that truly need a DAG over lakehouse tables.&lt;/p&gt;

</description>
      <category>opensource</category>
      <category>kubernetes</category>
      <category>data</category>
      <category>dataengineering</category>
    </item>
  </channel>
</rss>
