ApiaryActiveLive
Try: pause · settings · learn · wipe
← Community / Reading Room
AH
databases · 11 min read

Apache Hudi Incremental Pulls for Real‑Time Analytics

In a world where data moves at the speed of light, the ability to pull fresh insights in real time is no longer a luxury—it’s a necessity. For organizations…

In a world where data moves at the speed of light, the ability to pull fresh insights in real time is no longer a luxury—it’s a necessity. For organizations that rely on up‑to‑the‑minute analytics, the traditional batch‑processing paradigm simply cannot keep pace. Apache Hudi (Hadoop Upserts Deletes and Incrementals) tackles this challenge head‑on by providing a storage layer that supports ACID‑compliant upserts, deletes, and, most importantly, incremental pulls. Whether you’re a data engineer building a real‑time dashboard, a scientist tracking bee health metrics, or an AI agent orchestrating self‑governing workflows, Hudi’s incremental query model delivers the speed and consistency you need.

At its core, Hudi turns a data lake into a transactional engine. It stores data in Parquet or ORC files while maintaining a lightweight log of changes (the “delta log”). This log allows you to replay only the new or modified records since a given point in time—exactly what incremental pulls require. By combining Hudi’s efficient storage format with streaming engines like Spark Structured Streaming or Apache Flink, you can build pipelines that ingest, transform, and serve data with millisecond latency while preserving durability and consistency across distributed systems.

This pillar article dives deep into the mechanics of Hudi’s upsert, delete handling, and incremental query capabilities. We’ll walk through the underlying architecture, illustrate real‑world use cases (including a bee‑conservation data pipeline on the Apiary platform), and provide practical guidance on tuning performance for production workloads. By the end, you’ll have a clear understanding of how to harness Hudi for real‑time analytics and why it matters for data‑driven decision making.


Understanding Apache Hudi: A Primer

Apache Hudi (or “Hudi”) is an open‑source data lake framework designed to bring transactional semantics to large‑scale data lakes. Unlike traditional data lakes that treat storage as append‑only, Hudi introduces the concept of write‑ahead logs and commit metadata that enable ACID (Atomicity, Consistency, Isolation, Durability) guarantees on a distributed file system such as HDFS or Amazon S3.

Key components:

ComponentPurpose
Parquet/ORC filesDurable, columnar storage for committed data
Delta logJSON files that record every write operation (upsert, delete, compaction)
Commit metadataHolds commit timestamps, file lists, and partition information
CompactionBackground process that merges small delta files into larger ones for efficient reads

Hudi supports two main storage types:

  1. Copy‑On‑Write (COW) – Each write operation rewrites entire files. Provides faster reads but higher write cost.
  2. Merge‑On‑Read (MOR) – Writes are appended to delta logs, and data is merged on read. Offers lower write latency and is ideal for streaming workloads.

The choice between COW and MOR depends on your use case. For real‑time analytics where you need to read fresh data quickly, MOR is typically preferred. However, if your downstream systems require low‑latency point‑in‑time reads, COW may be more appropriate.


The Challenge of Streaming Data for Real‑Time Analytics

Streaming data pipelines often generate terabytes of records per day. Traditional batch jobs that run nightly or hourly cannot deliver the real‑time insights required by dashboards, alerting systems, or autonomous agents. The main pain points are:

  • Write‑Amplification: Re‑ingesting all data for each incremental run leads to unnecessary I/O and storage costs.
  • Consistency: Without transactional guarantees, concurrent reads can see partially written data, leading to stale or corrupted analytics.
  • Delete Propagation: Removing obsolete records in an append‑only system is problematic; stale data can accumulate and skew results.
  • Schema Evolution: Streaming sources frequently evolve their schema, and handling this gracefully is essential for long‑term maintainability.

Hudi addresses these challenges by treating the data lake as a mutable, versioned table. Every write is atomic, and the delta log records exactly what changed. Incremental pulls can then be expressed as “give me all rows that changed since commit X,” which is far more efficient than scanning the entire dataset.


Hudi’s Upsert Mechanism: How It Works hudi-upsert

An upsert (update or insert) is the cornerstone of Hudi’s data ingestion model. When a record arrives, Hudi determines whether it already exists in the table by looking up its primary key. If it does, Hudi updates the existing record; if not, it inserts a new one.

Primary Key and Partition Strategy

Hudi requires two pieces of metadata:

  1. Primary Key (hoodie.datasource.write.keygenerator.class) – A unique identifier for each row. In a bee‑conservation scenario, this might be a combination of hive ID and timestamp.
  2. Partition Path (hoodie.datasource.write.partitionpath.field) – A directory hierarchy that groups related rows. Partitioning by date (yyyy/MM/dd) is common for time‑series data.

Write Flow

  1. Ingest: A streaming source (e.g., Kafka) pushes a micro‑batch of records into Spark Structured Streaming.
  2. Generate Keys: Spark assigns primary keys and partition paths.
  3. Delta Log Append: Each micro‑batch writes a delta file (*.delta) containing the new and updated rows.
  4. Commit: Hudi writes a commit metadata file (commit-<timestamp>.json) that lists the delta files and updates the index.

Because the delta log is append‑only, writes are cheap: no need to rewrite large Parquet files. The index (often an embedded Bloom filter or parquet metadata) allows Hudi to locate the relevant delta files during a read without scanning the entire log.

Performance Highlights

  • Write Throughput: A single Hudi cluster on 10 nodes can ingest >200 MB/s of data with MOR storage, translating to >5 TB/day for a 50‑node cluster.
  • Latency: Micro‑batches of 1 second can be committed within 5–10 seconds, making near‑real‑time ingestion feasible.
  • Storage Efficiency: MOR writes typically produce 3–4x less storage overhead compared to naive append‑only writes because the index avoids duplicate data.

Handling Deletes in an Append‑Only World hudi-delete-handling

Deletes pose a unique problem in data lakes. Simply removing a file from S3 does not guarantee that downstream readers will not see stale data, because they may still reference a previously committed snapshot. Hudi solves this with a delete marker approach:

  1. Delete Marker: When a record needs to be deleted, Hudi writes a delta file containing only the primary key of the record and a special flag indicating deletion.
  2. Compaction: During compaction, Hudi merges the delete markers with the corresponding data files. The deleted records are physically removed from the Parquet/ORC files.
  3. Read Semantics: Hudi’s read engine consults the delta log and the latest index to filter out any records that have a delete marker.

This mechanism ensures that:

  • Consistency: Reads never return deleted records.
  • Efficiency: Deletion operations are lightweight because they only append a small delta file.
  • Recovery: The entire table can be rolled back to any prior commit, including before the delete, by using the commit timeline.

In the Apiary bee‑conservation pipeline, deletes are used to purge sensor data that has expired (e.g., older than 30 days). Because the delete marker is only a few bytes per record, the overhead is negligible even for millions of deletions per day.


Incremental Pulls: From Delta to Dashboard hudi-incremental-query

Incremental pulls are the heart of real‑time analytics. Hudi exposes a simple API for retrieving only the data that changed since a specific commit. This is achieved through the timeline client and the incremental query interface.

Incremental Query API

from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("HudiIncremental").getOrCreate()

hudi_path = "s3://apiary-bee-data/hudi/hive_table"

# Read the last commit timestamp from a metadata store
last_commit = "20240612T120000Z"

# Incremental read
df = spark.read.format("hudi") \
    .option("hoodie.datasource.query.type", "incremental") \
    .option("hoodie.datasource.query.start.instant", last_commit) \
    .load(hudi_path)

df.show()

The hoodie.datasource.query.start.instant parameter tells Hudi to read only the delta files committed after the specified instant. The API returns a DataFrame that contains only new and updated rows. If you also want to include deletes, set hoodie.datasource.query.end.instant accordingly.

Practical Use Cases

  • Real‑time Dashboards: A dashboard that refreshes every minute can call the incremental API to fetch only the latest 60 seconds of data, reducing latency dramatically.
  • Alerting Systems: An AI agent monitoring hive temperature can trigger alerts when new data indicates a sudden spike, without reprocessing the entire dataset.
  • Machine Learning Pipelines: Incremental pulls can feed a streaming model that updates predictions on the fly, such as estimating hive health scores.

Performance Considerations

MetricTypical Value (MOR)Typical Value (COW)
Read Latency (per micro‑batch)1–3 seconds5–10 seconds
Storage Overhead (delta log)5–10% of data size2–5%
Query Throughput200 MB/s80 MB/s

The incremental API is designed to work seamlessly with Spark Structured Streaming, Flink, or any engine that can consume Parquet/ORC. Because the delta log is sorted by commit instant, reads are highly cacheable and can be pipelined with downstream transformations.


Integrating Hudi with Spark Streaming and Flink spark-streaming flink-integration

Spark Structured Streaming

Spark’s integration with Hudi is mature and straightforward. The typical pattern is:

  1. Source: Kafka or Kinesis streams.
  2. Transform: User logic (e.g., normalizing timestamps, computing derived metrics).
  3. Sink: Hudi table with write.format("hudi").
val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "kafka:9092")
  .option("subscribe", "bee_hive_events")
  .load()
  .selectExpr("CAST(value AS STRING)")
  .select(from_json(col("value"), schema).as("data"))
  .select("data.*")

df.writeStream
  .format("hudi")
  .option("hoodie.table.name", "bee_hive_events")
  .option("hoodie.datasource.write.operation", "upsert")
  .option("hoodie.datasource.write.keygenerator.class", "org.apache.hudi.keygen.SimpleKeyGenerator")
  .option("hoodie.datasource.write.partitionpath.field", "date")
  .option("hoodie.datasource.write.recordkey.field", "hive_id")
  .option("checkpointLocation", "/tmp/checkpoints/bee_hive_events")
  .outputMode("append")
  .start()

The checkpointLocation stores the offset for Kafka, ensuring at-least-once semantics. Hudi’s upsert operation guarantees that duplicate hive events are merged correctly.

Flink Integration

Flink users can leverage the Hudi connector, which supports both streaming and batch modes. The Flink API is similar to Spark’s but uses the HoodieSink and HoodieSource classes.

DataStream<JSONObject> source = env
  .addSource(new FlinkKafkaConsumer<>("bee_hive_events", new SimpleStringSchema(), props))
  .map(JSON::parseObject);

HoodieSink<String> hudiSink = HoodieSink.builder()
  .withTableName("bee_hive_events")
  .withPrimaryKey("hive_id")
  .withPartitionField("date")
  .withOperation(HoodieOperation.UPSERT)
  .withPath("s3://apiary-bee-data/hudi/bee_hive_events")
  .build();

source.addSink(hudiSink);

Flink’s low‑latency streaming model, combined with Hudi’s incremental query, makes it ideal for scenarios where data must be processed in micro‑seconds, such as real‑time anomaly detection in hive temperature.

Performance Tuning Tips

Tuning ParameterEffectRecommended Value
hoodie.parquet.block.sizeLarger blocks reduce overhead128 MB
hoodie.compact.inlineInline compaction for small tablestrue
hoodie.compact.inline.max.delta.commitsControls when inline compaction triggers5
spark.default.parallelismNumber of shuffle partitions2 × number of cores
spark.sql.shuffle.partitionsShuffle size200–400

By fine‑tuning these parameters, you can balance write throughput, read latency, and storage efficiency to match your workload.


Case Study: Bee Conservation Data Pipeline

Background

The Apiary platform monitors thousands of beehives worldwide. Each hive is equipped with temperature, humidity, and vibration sensors that emit data every second. Conservationists need near‑real‑time insights to detect stressors such as heat waves or predator attacks.

Pipeline Overview

  1. Ingestion: Sensors publish events to an MQTT broker, which forwards them to Kafka.
  2. Streaming Layer: Spark Structured Streaming reads from Kafka, performs basic cleaning, and writes to Hudi (MOR).
  3. Incremental Pull: A Flink job reads the Hudi incremental stream every 30 seconds, aggregates hive health scores, and writes the results to a downstream ClickHouse cluster for visualization.
  4. Alerting: A rule‑based AI agent monitors the ClickHouse table. When a hive’s health score drops below a threshold, the agent triggers an SMS alert to the local beekeeper.

Numbers

MetricValue
Data Volume5 TB/day (≈ 500 million events)
Write Throughput250 MB/s (peak)
Latency15 seconds from sensor to alert
Storage Cost30% reduction compared to append‑only S3 storage

The key to achieving these numbers was Hudi’s incremental pull. Instead of re‑ingesting the entire 5 TB, the Flink job fetched only the new 30‑second slice (~1.2 GB). This drastically reduced compute time and network usage, enabling the system to scale to 10,000 hives without additional infrastructure.

Lessons Learned

  • Partition by Hive ID: While date partitioning is common, in this scenario partitioning by hive ID improved locality for per‑hive analytics.
  • Delta Commit Size: Setting hoodie.delta.commits to 3 seconds ensured that each micro‑batch fit comfortably into a single delta file, simplifying compaction.
  • Compaction Strategy: Using background compaction with a 24‑hour window kept read performance high without impacting writes.

Best Practices & Performance Tuning

CategoryRecommendationRationale
Schema DesignKeep schema flat; avoid deeply nested structuresParquet/ORC performance drops with high nesting
PartitioningUse high‑cardinality fields (e.g., hive ID, date)Reduces scan size for targeted queries
CompactionEnable hoodie.compact.inline=true for small tablesEliminates background jobs
IndexingEnable hoodie.datasource.write.precombine.fieldImproves upsert performance by reducing duplicate rows
RetentionUse hoodie.retention.check.enabled=trueKeeps delta log size manageable
MonitoringExpose Hudi metrics to PrometheusEnables proactive scaling

Example: Enabling Precombine

hoodie.datasource.write.precombine.field: "timestamp"

When multiple updates for the same primary key arrive, Hudi uses the precombine field to keep only the newest record in the delta file. This reduces the size of delta files and speeds up compaction.

Garbage Collection & Cleanup

Hudi’s hoodie.retention.check.interval.ms controls how often the system removes obsolete files. Setting this to 1 hour for high‑velocity pipelines ensures that storage does not balloon.


Future Directions: Hudi and AI Agents

As AI agents evolve from reactive scripts to self‑governing entities, they will need to ingest, process, and respond to data in real time. Hudi’s incremental pull model aligns perfectly with this vision:

  • Autonomous Decision Making: An agent can query the latest data snapshot every few seconds, evaluate conditions, and trigger actions (e.g., opening hive doors, adjusting ventilation).
  • Self‑Healing Pipelines: The agent can monitor Hudi’s commit timeline, detect stalled commits, and automatically retry or rollback.
  • Policy Enforcement: Using Hudi’s commit hooks, agents can enforce data governance policies (e.g., GDPR “right to be forgotten”) by programmatically issuing deletes.

In the context of bee conservation, imagine an AI agent that monitors hive temperature in real time, predicts a heat‑wave event, and automatically orders a cooling unit to be dispatched, all without human intervention. Hudi provides the reliable, incremental data foundation that makes such autonomy feasible.


Why it Matters

Apache Hudi’s incremental pull capability transforms how we approach real‑time analytics on data lakes. By enabling efficient upserts, robust delete handling, and precise incremental queries, Hudi bridges the gap between the high‑throughput world of streaming data and the analytical demands of modern enterprises. Whether you’re protecting pollinators, orchestrating autonomous AI agents, or powering a global business intelligence platform, Hudi gives you the speed, consistency, and scalability you need to act on data as soon as it arrives.


Frequently asked
What is Apache Hudi Incremental Pulls for Real‑Time Analytics about?
In a world where data moves at the speed of light, the ability to pull fresh insights in real time is no longer a luxury—it’s a necessity. For organizations…
What should you know about understanding Apache Hudi: A Primer?
Apache Hudi (or “Hudi”) is an open‑source data lake framework designed to bring transactional semantics to large‑scale data lakes. Unlike traditional data lakes that treat storage as append‑only, Hudi introduces the concept of write‑ahead logs and commit metadata that enable ACID (Atomicity, Consistency, Isolation,…
What should you know about the Challenge of Streaming Data for Real‑Time Analytics?
Streaming data pipelines often generate terabytes of records per day. Traditional batch jobs that run nightly or hourly cannot deliver the real‑time insights required by dashboards, alerting systems, or autonomous agents. The main pain points are:
What should you know about hudi’s Upsert Mechanism: How It Works hudi-upsert?
An upsert (update or insert) is the cornerstone of Hudi’s data ingestion model. When a record arrives, Hudi determines whether it already exists in the table by looking up its primary key. If it does, Hudi updates the existing record; if not, it inserts a new one.
What should you know about write Flow?
Because the delta log is append‑only, writes are cheap: no need to rewrite large Parquet files. The index (often an embedded Bloom filter or parquet metadata ) allows Hudi to locate the relevant delta files during a read without scanning the entire log.
References & sources
  1. Apiary Reading Room — Open, cited knowledge base — funded to keep bee & practical research free.
From the Apiary Reading Room. Opinion & editorial — not financial advice. We don't overclaim.
More from the Reading Room