Hive was built for batch: land files in HDFS or object storage, run a load or an INSERT ... SELECT, and query the result an hour later. Many pipelines still need fresher data than that. Clickstream, IoT readings and application logs arrive continuously, and analysts want them queryable within minutes, not after a nightly job. Copying small files into a table directory every few seconds is the obvious hack, and it fails in two ways: readers can see half-written files, and the table drowns in tiny files that make every later query slow.

Hive's Streaming Data Ingest API is the supported answer. A client opens a connection to a transactional ORC table, writes records inside Hive ACID transactions, and commits them so that readers see either all of a commit or none of it. The compactor later merges the many small delta files into large ones. This article explains the mechanism from first principles, shows a complete Java client for the V2 API that ships with Hive 3 and later, works through the commit cadence arithmetic for a realistic feed, and lists what breaks in production. Every API name below was checked against the Apache Hive wiki page for Streaming Data Ingest V2 on 2026-10-03. For the ACID machinery underneath, read Hive ACID transactions alongside this page.

How a streaming commit becomes visible

Hive ACID tables do not modify files in place. Each transaction that writes to a table is allocated a monotonically increasing write ID for that table by the metastore, and its rows land in a new directory named after the range of write IDs it covers, such as delta_0000102_0000102. A reader starts by asking the metastore for a snapshot: the highest write ID it may see, plus the list of write IDs that are open or aborted below it. It then reads the base directory and every delta inside that snapshot and ignores anything else. That is why a commit is atomic: the files may already be on disk, but no reader trusts them until the metastore records the transaction as committed.

Streaming ingest is a thin client on top of this. The client asks the metastore to open a transaction, writes ORC rows into a delta directory for its write ID, and then commits. If the client dies, the transaction stops heartbeating, the metastore times it out and marks it aborted, and readers skip its rows. Nothing has to be cleaned up by hand before the next reader arrives.

The cost of this design is files. Every commit produces at least one new delta per bucket it touched, in every partition it touched. Reads must merge base plus all deltas, so a table with thousands of deltas is slow to read and expensive for the NameNode. The compactor fixes this in the background. A minor compaction merges deltas into one larger delta; a major compaction rewrites base plus deltas into a new base. The cleaner removes obsolete directories once no reader still needs them. Streaming ingestion only works if compaction keeps up, and most operational advice in this article follows from that one fact. The compaction article covers the compactor itself in depth.

Agent 1..NJava client, one conn eachRecord writerStrict delimited / JSON / regexHive Metastoretxn + write IDs, locks, heartbeatsopen/commit txnTable dir (ORC, ACID)base_0000100/delta_0000101_0000101/delta_0000102_0000102/delta_0000103_0000110/ ...ORC deltasCompactorinitiator + workers + cleanerminor / majorthresholdsReadersHive, Impala, Sparkvalid write-ID listA commit makes a delta visible atomically; the compactor merges deltas so reads stay fast; aborted write IDs are skipped.
Agents write ORC deltas inside metastore transactions; readers trust only committed write IDs; the compactor merges deltas in the background.

Preparing the table and the cluster

The V2 API has three hard requirements on the target table: it must be transactional, it must be stored as ORC, and the cluster must run the DbTxnManager. Bucketing is supported but no longer required, which was a real pain point of the original API. Partitioning is optional, and if you do not supply partition values the connection can route rows to partitions dynamically.

-- Target table: transactional ORC, partitioned by event date and hour.
CREATE TABLE web.clicks (
  event_id    STRING,
  user_id     STRING,
  url         STRING,
  referrer    STRING,
  ts          TIMESTAMP
)
PARTITIONED BY (dt STRING, hr STRING)
STORED AS ORC
TBLPROPERTIES ('transactional' = 'true');

On the cluster side, the settings the wiki lists are the transaction manager and a working compactor. Without worker threads, compaction requests queue up and nothing ever merges, which is the single most common reason a streaming table degrades over weeks.

<!-- hive-site.xml on the metastore / HiveServer2 hosts -->
<property><name>hive.support.concurrency</name><value>true</value></property>
<property><name>hive.txn.manager</name>
          <value>org.apache.hadoop.hive.ql.lockmgr.DbTxnManager</value></property>
<property><name>hive.compactor.initiator.on</name><value>true</value></property>
<property><name>hive.compactor.cleaner.on</name><value>true</value></property> <!-- Hive 4+ -->
<property><name>hive.compactor.worker.threads</name><value>4</value></property>

Run the initiator and cleaner in exactly one metastore instance, and give worker threads to hosts with spare CPU and memory. Compaction jobs are real MapReduce or Tez jobs and they compete with user queries, so put them in their own YARN queue with hive.compactor.job.queue and size that queue for your peak ingest rate, not your average.

Writing a streaming client

A streaming client has a simple lifecycle: build a record writer, build a connection, then loop over begin, write many records, commit. The record writer converts your bytes into table rows. StrictDelimitedInputWriter expects delimited text whose fields match the table columns exactly, StrictJsonWriter expects JSON objects whose keys match column names, and StrictRegexWriter parses text with a regular expression. The word Strict matters: a record that does not match the schema fails the write rather than being silently coerced.

import org.apache.hadoop.hive.conf.HiveConf;
import org.apache.hive.streaming.HiveStreamingConnection;
import org.apache.hive.streaming.StrictJsonWriter;
import org.apache.hive.streaming.StreamingException;

public class ClickStreamer implements AutoCloseable {
    private final HiveStreamingConnection conn;
    private int pending = 0;
    private long openedAt;

    public ClickStreamer(HiveConf conf, String agentId) throws StreamingException {
        StrictJsonWriter writer = StrictJsonWriter.newBuilder().build();
        // No static partition values: rows carry dt and hr as their last two fields,
        // so the connection routes each row to its partition (dynamic mode).
        this.conn = HiveStreamingConnection.newBuilder()
                .withDatabase("web")
                .withTable("clicks")
                .withAgentInfo(agentId)          // shows up in metastore txn records
                .withRecordWriter(writer)
                .withHiveConf(conf)
                .connect();
        begin();
    }

    private void begin() throws StreamingException {
        conn.beginTransaction();
        pending = 0;
        openedAt = System.currentTimeMillis();
    }

    /** Called for every event; commits on size or age, whichever comes first. */
    public void write(byte[] jsonRecord) throws StreamingException {
        conn.write(jsonRecord);
        pending++;
        if (pending >= 200_000 || System.currentTimeMillis() - openedAt >= 60_000) {
            conn.commitTransaction();   // rows become visible to new readers here
            begin();
        }
    }

    public void fail() {
        try { conn.abortTransaction(); } catch (StreamingException ignored) { }
    }

    @Override public void close() throws StreamingException {
        if (pending > 0) conn.commitTransaction();
        conn.close();
    }
}

Three details in this client matter. First, it commits on size or age, so a quiet period never leaves rows invisible for long and a burst never builds an enormous transaction. Second, the upstream offset (a Kafka offset, a file position) must be advanced only after commitTransaction() returns; advancing it earlier turns a crash into data loss. Third, the connection runs its own heartbeat thread at half of hive.txn.timeout, so a slow upstream does not abort the open transaction, but a JVM that is frozen in a long garbage-collection pause can still miss heartbeats.

Exactly-once delivery is not free. If the process crashes after the Hive commit but before the offset is saved, the replay writes the same rows again. Either store the upstream offset in a place you update atomically with the data, or make rows idempotent with a stable event_id and deduplicate with a MERGE or a windowed query downstream.

Commit cadence and the delta bill

Commit cadence is the main tuning decision, and it is a direct trade between freshness and file count. The wiki's guidance is to group thousands of records per transaction rather than commit every record, and its description of transaction batches explains why: transactions in one batch write to the same physical file per bucket, so batching reduces the number of files and delta directories. Smaller commits mean fresher data and more deltas; larger commits mean fewer deltas and staler data.

Deltas written with streaming optimizations are also cheaper to write and more expensive to read than normal ORC: the wiki notes that dictionary encoding, indexes and compression are disabled for write throughput, and that the compactor rewrites them into the optimised format. So until compaction runs, recent data is both fragmented and less compact. That is acceptable for a few hours; it is not acceptable for weeks.

Commit intervalFreshnessDeltas per agent per partition per hourTypical use
10 sSeconds360Avoid unless compaction is very aggressive
60 sAbout a minute60Dashboards with minute-level freshness
5 minMinutes12Most analytics feeds
15 minQuarter hour4Large tables where read cost dominates

The compactor's triggers are thresholds on what accumulates. hive.compactor.delta.num.threshold starts a minor compaction once a partition has more than that many deltas, hive.compactor.delta.pct.threshold starts a major compaction once deltas reach that fraction of the base size, and hive.compactor.abortedtxn.threshold reacts to piles of aborted transactions. Check the defaults on your distribution rather than assuming them; what matters is that the thresholds and the worker count together can keep up with the rate your commits produce deltas.

Worked example: a clickstream feed

Take a clickstream feed of 20,000 events per second, about 400 bytes of JSON each, consumed from Kafka by eight agents, written into web.clicks partitioned by date and hour. Analysts want data within five minutes.

  1. Throughput per agent. 20,000 / 8 = 2,500 events per second per agent, about 1 MB/s of input each. One JVM handles that comfortably.
  2. Pick the commit rule. A 60-second age limit gives one-minute freshness, comfortably inside the five-minute target. In 60 seconds an agent sees 150,000 events, so the 200,000-record size limit in the client above only triggers during bursts.
  3. Count deltas. Each agent writes to the current hour's partition, so each partition receives 8 agents x 60 commits = 480 deltas in its hour, plus a trickle of late events afterwards. Without compaction, a query over one day reads more than 11,000 delta directories.
  4. Size compaction. With a delta-count threshold in the tens, the current partition is minor-compacted several times an hour. An hour holds 20,000 x 3,600 x 400 bytes, about 29 GB of raw JSON and much less as ORC, and each compaction handles only the slice since the last, so start with a few worker threads in a dedicated queue and measure. Schedule a major compaction for each hour partition once it is closed, so that yesterday's data is a single optimised base.
  5. Check the result. Run SHOW COMPACTIONS; and SHOW TRANSACTIONS; for a day. Healthy means compactions move from initiated to succeeded within minutes and open transactions are only the eight agents' current ones.

At ten-second freshness the same arithmetic gives 2,880 deltas per partition per hour; use a different tool.

Failure modes

SymptomCauseFix
Queries get slower every dayCompactor not running: no worker threads, initiator off, or the compaction queue starvedCheck SHOW COMPACTIONS; enable workers; give compaction its own queue
Thousands of tiny files per partitionCommitting every few records or every few secondsCommit on size and age; target tens of seconds to minutes
Transactions aborted by timeoutAgent GC pauses or network partitions longer than hive.txn.timeoutTune heap; alert on aborts; keep the default heartbeat logic
Write fails with a schema errorStrict writer saw a record whose fields do not match the tableValidate upstream; route bad records to a dead-letter topic before the writer
Duplicate rows after a restartOffset saved after crash point, data already committedSave offsets after commit and deduplicate on a stable event ID
Readers see old data onlySnapshot taken before the commit, or a long-running queryExpected: snapshots are fixed per query; rerun the query
Many aborted transactions accumulateCrash loops in agentsFix the loop; the aborted-transaction threshold triggers cleanup compaction
Too many partitions open at onceDynamic mode with late or skewed events touching many partitions per commitBound lateness, route very late data to a batch path

The most dangerous failure is silent: everything works for a month, then compaction falls behind after a traffic increase, and read latency climbs a little each day until someone notices. Alert on deltas per partition and on the age of the oldest initiated compaction, not only on agent liveness.

The small files article explains why the NameNode and object stores both suffer from many small files, and why the cost shows up in places that look unrelated to ingestion.

Trade-offs and alternatives

Streaming into Hive ACID is a good fit when the table already lives in Hive, consumers query it with Hive, Impala or Spark through the metastore, and freshness of one to fifteen minutes is enough. It is a poor fit for second-level latency, for very high update rates, or when you need readers outside the Hive ecosystem.

ApproachFreshnessStrengthsCosts
Streaming Ingest V2 into ACID ORCMinutesAtomic commits, no external engine, Hive-nativeDelta files, depends on compactor health, ORC only
Micro-batch INSERT ... SELECT every N minutesN minutesSimple, uses ordinary queriesQuery startup cost per batch; still produces deltas on ACID tables
Land files, add partitions hourlyHourNo transactions, cheapest readsStale; readers may see partial files unless you write to a staging path and move
Spark or Flink into an Iceberg tableMinutesOpen table format, snapshot isolation, many enginesSeparate snapshot expiry and file compaction to operate

Many teams moving off Hive-only stacks choose the last option; the Hive and Iceberg article explains how Hive reads Iceberg tables so you can keep existing queries while changing the ingest path.

What to do next

  1. Confirm your Hive version is 3.0 or later and that the table is transactional ORC; check DESCRIBE FORMATTED for transactional=true.
  2. Enable the DbTxnManager, the compaction initiator and cleaner on one metastore, and give worker threads a dedicated YARN queue.
  3. Build the client with a size-and-age commit rule; start at 60 seconds and measure deltas per partition per hour.
  4. Advance upstream offsets only after commitTransaction() returns, and carry a stable event ID for deduplication.
  5. Put a schema check and dead-letter route in front of the strict record writer.
  6. Alert on deltas per partition, the age of the oldest pending compaction and the aborted transaction count.
  7. Schedule a major compaction for closed partitions and verify that query time on yesterday's data matches a batch-loaded table.
  8. If you need second-level freshness, stop and evaluate a different engine before tuning further.
Key takeaway: Hive streaming ingestion writes records into transactional ORC tables inside metastore transactions, so each commit becomes visible atomically and a crashed writer's rows are skipped. Every commit also creates delta files, so the design is a trade between freshness and file count, and it only works if the compactor keeps up. Commit on size and age, advance upstream offsets only after a commit, give compaction its own resources and alert on deltas per partition.