Skip to content

Kafka to Iceberg: too many small files, and how to compact them

Iceberg
Chad Harris·October 3, 2026·11 min read

Query planning on an Iceberg table fed from Kafka keeps getting slower, the files metadata table lists thousands of data files of a few kilobytes each, and the metadata directory keeps growing even though the data volume is steady.

The small files problem is a table whose data is spread over many files far below the target file size, because a streaming writer commits a new set of files on every checkpoint or commit interval. The usual fix has three parts: compact the partitions the writer has finished, make each commit add fewer files by lengthening the commit interval or sending each partition to fewer writers, and expire snapshots while keeping the last snapshot the Flink job created.

The Iceberg maintenance documentation states the cost directly. More data files means more metadata stored in manifest files, and small data files cause an unnecessary amount of metadata and less efficient queries from file open costs. The same page says streaming queries may produce small data files that should be compacted into larger files.

The cost is shared. Iceberg lets engines such as Spark, Trino, Flink, Presto, Hive and Impala work with the same tables at the same time, as the Apache Iceberg home page puts it, and Snowflake documents Iceberg tables that sit in external cloud storage the owner manages. A table with a poor file layout therefore costs every engine that reads it, not only the one that writes it. For how tables, snapshots and metadata fit together, see the complete Iceberg guide.

This page is about the Iceberg table a streaming job writes into. If the symptom is a Kafka consumer group falling behind and no Iceberg sink is involved, the guide to Kafka consumer lag is the better starting point. If the Flink job writing the table restarts because its checkpoints fail, start with Flink checkpoints that fail or time out.

Why a streaming sink writes small files

A streaming sink commits on a schedule, and each commit closes the files that are open at that moment, whatever their size. Three things set how many files a commit adds: how often the sink commits, how many writer tasks run, and how many partitions each writer touches between commits.

The Flink writes documentation says the Iceberg commit happens after a successful Flink checkpoint, in the notifyCheckpointComplete callback, and its alerting example treats the checkpoint interval as the expected Iceberg commit interval. Writer subtasks flush and upload their files during the checkpoint. So the checkpoint interval is the commit interval, and every checkpoint adds a new set of files.

execution.checkpointing.interval

The Flink checkpointing documentation describes this setting as the base interval at which checkpoints are scheduled.

Kafka Connect commits on a timer

The Apache Iceberg Sink Connector for Kafka Connect is the open source, self-managed way to write Kafka topics into Iceberg tables, installed by copying its distribution archive into the plugins directory of a Kafka Connect cluster the operator runs. It coordinates commits centrally through a control topic. Its commit interval defaults to 300,000 ms, five minutes.

iceberg.control.commit.interval-ms=300000

A failing connector task stops commits altogether, which is a different problem from small files. The guides to diagnosing a failed Kafka Connect connector and to what Kafka Connect is cover that case.

Writers and partitions multiply the count

Each writer task writes its own files for each partition it receives data for. The Flink configuration documentation gives an example of a table partitioned by day, where one writer task receives the oldest 150 days of a long tail and writes 150 small files, one per day, and notes that flushing and uploading those 150 files at checkpoint time can be slow. A short interval, a high writer parallelism and a fine-grained partition spec all raise the file count per commit, and they multiply together.

Measure the problem before changing anything

Pick the metric first, then change one setting at a time and watch that metric.

File sizes and counts from the metadata tables

Iceberg exposes a table’s current files and partitions as metadata tables, which the Spark queries documentation lists. The maintenance documentation names the files metadata table as the place to inspect data file sizes and decide when to compact a partition.

SELECT partition, count(*) AS files, avg(file_size_in_bytes) AS avg_bytes
FROM prod.db.table.files
GROUP BY partition
ORDER BY files DESC;

The partitions metadata table gives the same count without reading every file entry.

SELECT partition, file_count, total_data_file_size_in_bytes
FROM prod.db.table.partitions;

Compare avg_bytes with the table’s target file size, which defaults to 536,870,912 bytes (512 MB) in the table configuration reference.

write.target-file-size-bytes

Files per commit from the snapshots table

Each snapshot’s summary records what the commit added, including added-data-files and total-data-files. Reading the snapshots table over a day shows how many files each commit brings in and how fast the total grows.

SELECT committed_at, summary['added-data-files'], summary['total-data-files']
FROM prod.db.table.snapshots
ORDER BY committed_at DESC;

The Flink Iceberg sink publishes writer metrics under IcebergStreamWriter and committer metrics under IcebergFilesCommitter, listed in the Flink writes documentation.

Metric Emitted by What it shows
flushedDataFiles IcebergStreamWriter Data files flushed by the writers
dataFilesSizeHistogram IcebergStreamWriter The size distribution of those files
lastFlushDurationMs IcebergStreamWriter How long the last flush took
committedDataFilesCount IcebergFilesCommitter Data files committed to the table
elapsedSecondsSinceLastSuccessfulCommit IcebergFilesCommitter Seconds since the last successful commit

The histogram metrics need org.apache.flink:flink-metrics-dropwizard on the classpath, which Flink does not ship by default. The documentation calls elapsedSecondsSinceLastSuccessfulCommit an ideal alerting metric for failed or missing commits, and gives the example of a five-minute checkpoint interval with an alert at more than 60 minutes.

Connector task errors and consumer lag

Small files are a symptom on the table side. Two signals on the Kafka side show whether ingestion into the table keeps up: connector task errors, for example a task failing after a backward-incompatible schema change, and consumer lag, which shows ingestion falling behind. The Kafka Connect user guide describes a REST endpoint that returns the status of a connector, including error information and the state of every task.

GET /connectors/{name}/status

Consumer lag has no single definition. The consumer metric records-lag-max is, in the Kafka monitoring documentation, based on the current offset and not the committed offset, while Burrow, an open source lag checker, monitors committed offsets for all consumers. Two dashboards can therefore disagree about the same group, so check which definition each one uses before comparing them.

Lag in records is also an estimate. The Kafka design documentation notes that after log compaction every offset remains a valid position even when its record has been compacted away, and KIP-98 introduced control messages that transactions write into topics. A gap between offsets is therefore an upper bound on the records still to be read, not an exact count.

Compact the files that are already there

Compaction rewrites existing small files into larger ones. The Spark procedures documentation describes rewrite_data_files as combining small files into larger files to reduce metadata overhead and runtime file open cost.

CALL catalog_name.system.rewrite_data_files('db.sample');

With no options the procedure uses the binpack strategy, which combines small files and splits large ones according to the table’s write size. These defaults, from the Spark procedures documentation, decide which files it picks.

Option Default What it decides
target-file-size-bytes 536870912 (512 MB), from write.target-file-size-bytes The size the output files aim for
min-file-size-bytes 75% of the target Files under this size are rewritten regardless of other criteria
max-file-size-bytes 180% of the target Files over this size are rewritten regardless of other criteria
min-input-files 5 A file group with this many files or more is rewritten regardless of other criteria
max-concurrent-file-group-rewrites 5 How many file groups are rewritten at the same time
partial-progress.enabled false Whether groups of files are committed before the whole rewrite completes

Limit the rewrite to closed partitions

A streaming writer keeps committing to the newest partitions. The where argument selects only files that may contain matching data, so a rewrite can target partitions the sink has finished writing.

CALL catalog_name.system.rewrite_data_files(
  table => 'db.sample',
  where => 'event_date < "2026-10-01"'
);

Commit large rewrites in pieces

On a large table, one commit at the end of a long rewrite is a risk, because a failure late in the job loses the work. Partial progress commits groups of rewritten files as they finish, and the number of commits it may produce is capped at 10 by default.

CALL catalog_name.system.rewrite_data_files(
  table => 'db.sample',
  options => map('partial-progress.enabled', 'true')
);

With partial progress on, the procedure output includes a count of files that failed to rewrite.

failed_data_files_count

Compaction can also run on a schedule. The AWS Prescriptive Guidance on Iceberg compaction lists scheduled bin packing as its own use case, using the Athena OPTIMIZE statement when the number of small files is unknown, or Amazon EMR or AWS Glue when large volumes are expected. It advises running OPTIMIZE per table partition to avoid timeouts, and a scheduled job per table follows the same per-table, per-partition scoping as the where argument above.

Row-level changes add delete files

A table that receives updates or deletes, for example from change data capture, has a second source of small files. Iceberg defines the table properties write.delete.mode, write.update.mode and write.merge.mode, each set to copy-on-write or merge-on-read, in its table configuration reference, and the default for all three is copy-on-write. The AWS Prescriptive Guidance on writing to Iceberg tables explains the trade. Merge-on-read speeds up writes because updates and deletes are stored as separate small files, which readers merge with the base files, and it is typically suited to streaming workloads with updates. It also requires regular compaction so read performance does not degrade over time.

Teams without Spark can compact from Flink. The Flink maintenance documentation describes a rewrite files action that behaves the same as Spark’s rewriteDataFiles, and a TableMaintenance API that schedules RewriteDataFiles, for example after a set number of new data files, with a target size and partial progress. It needs a lock so that two maintenance tasks never run on the same table at once, either through JDBC, ZooKeeper or a Flink-maintained lock when no parallel maintenance job runs for that table.

The Flink writes documentation also describes IcebergSink, the SinkV2 based sink, which can run rewriteDataFiles() after each commit so no separate job is needed. The same page calls the SinkV2 implementation an experimental feature to use with caution.

F1 Compaction and sink settings, side by side
Compaction (rewrite_data_files) Sink settings
What it changes Files already in the table, rewritten into larger ones. How many files each new commit adds.
Metric that proves it worked file_count per partition falls, and file_size_in_bytes moves toward the target. added-data-files per snapshot and flushedDataFiles per checkpoint fall.
Main settings target-file-size-bytes, min-input-files, partial-progress.enabled, where. Checkpoint or commit interval, distribution-mode, write-parallelism.
Cost A separate job that reads and rewrites data, and adds a commit. Longer intervals delay data, and hash distribution adds a shuffle.
Drawn from the Apache Iceberg documentation on table maintenance, Spark procedures, Flink writes, Flink configuration and the Kafka Connect sink. Iceberg maintenance documentation

Change the sink so fewer small files arrive

Compaction cleans up after the writer. Sink settings reduce how many files the writer produces. Change one at a time and compare added-data-files per snapshot before and after.

Lengthen the commit interval

A longer interval gives each writer more data per file. For Flink that means a longer checkpoint interval, set by the first key below, and for Kafka Connect a larger commit interval, set by the second.

execution.checkpointing.interval
iceberg.control.commit.interval-ms

The cost is that data reaches the table less often. For Flink, a longer interval also changes recovery, because the job restores from its last completed checkpoint, so read the Flink checkpoint guide before stretching it.

Route each partition to fewer writers

Hash distribution shuffles rows by partition key before the writers, so each partition goes to one writer instead of all of them. The table property below sets it for the table, and the Flink distribution-mode write option overrides it for one write.

write.distribution-mode=hash

The Flink writes documentation lists its limits. It does not handle skewed data well, and writer parallelism is capped at the cardinality of the hash key, so with 10 distinct keys at most 10 writer tasks get traffic. Range distribution, which the same page marks as experimental, collects traffic statistics to balance skewed partitions. The shuffle itself adds CPU for serialisation and network I/O.

Match writer parallelism to the data rate

The Flink write-parallelism option overrides the writer parallelism, which otherwise follows the upstream operator. Fewer writers each write larger files. Watch lastFlushDurationMs and checkpoint duration after lowering it, because fewer writers each flush more data at checkpoint time.

write-parallelism

Revisit the partition spec

A partition spec finer than the data rate gives each commit many partitions with a few rows each. The Iceberg partitioning documentation says Iceberg can partition timestamps by year, month, day and hour, and an hourly spec on a low-volume topic is the case where a writer sees only a few rows per partition between two commits. A file cannot grow past the data its writer receives for that partition between commits, whatever the target size is. The Flink write option below overrides the table’s target for one write, as the Flink configuration documentation lists, and raising it does not change that limit.

target-file-size-bytes

Clean up snapshots and metadata without breaking the writer

Every commit also writes a new metadata file and a new snapshot. The maintenance documentation says tables with frequent commits, like those written by streaming jobs, may need to clean metadata files regularly, and recommends expiring snapshots to delete data files that are no longer needed and keep table metadata small.

The first table property below makes the writer delete the oldest tracked metadata file on each commit, and the second sets how many previous metadata files it keeps, 100 by default.

write.metadata.delete-after-commit.enabled=true
write.metadata.previous-versions-max=100

Snapshots are expired with the expire_snapshots procedure or the expireSnapshots action. Expired snapshots are no longer available for time travel. Compaction leaves the old small files referenced by older snapshots until those snapshots expire, so storage only shrinks after expiry.

CALL catalog_name.system.expire_snapshots(
  table => 'db.sample',
  older_than => TIMESTAMP '2026-10-01 00:00:00'
);

The Flink writes documentation adds a warning for Flink sinks. The job keeps its last committed checkpoint ID in the snapshot summary and stores uncommitted data as temporary files, so expiring snapshots and deleting orphan files can corrupt the Flink job’s state. Before expiring on a table a Flink job writes to, work through this order.

  1. Find the last snapshot the Flink job created, which carries the flink.job-id property in its summary.
  2. Choose an expiry cutoff that keeps that snapshot. A compaction commit made after it is newer, so the newest snapshot is not necessarily the one the job needs.
  3. Keep orphan file retention at or above the longest write. The maintenance documentation sets the default at 3 days and warns that a shorter interval than the longest write can corrupt the table.
  4. Run the expiry.

Where Kpow and Flex fit

This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, and Kpow, its Kafka management and monitoring product, so weigh this section with that in mind. Neither product’s documentation describes Iceberg table file counts, metadata tables or compaction, so file sizes and compaction stay in the Iceberg metadata tables and the Spark or Flink maintenance jobs above.

Kafka side. The Factor House talk Journey to an open lakehouse names two signals to surface while Kafka feeds Iceberg: connector task errors, such as a failure after a backward-incompatible schema change, and Kafka consumption lag, which shows ingestion falling behind. Kpow documents both. Its Kafka Connect management documentation describes pausing, restarting and editing connectors, and restarting individual tasks or viewing their stack traces when a task is in an ERROR state. Its Prometheus metrics glossary lists consumer group lag metrics such as group_offset_lag.

Flink side. Flex documents these views for each job in its job inspection documentation.

  • The Checkpoints tab shows checkpoint counts by status, average checkpoint size and duration, and a history table of every attempt.
  • The Configuration view shows the configuration the job was submitted with, including its parallelism and a map of user-provided values.
  • The Events tab is a reverse-chronological lifecycle log, so a failed checkpoint and a restart sit in one place.
  • Subtask metrics show per-instance figures such as Bytes received and Records sent.

The Iceberg sink commits after each successful checkpoint, as the Flink writes documentation states, so the checkpoint history shows when commits could have happened. It does not show whether an Iceberg commit succeeded, because the same documentation notes that commits can fail while checkpoints keep succeeding.

The job documentation does not describe the Iceberg sink metrics such as flushedDataFiles or elapsedSecondsSinceLastSuccessfulCommit, so those go through Flink’s metric reporters to an external system. Flex is one of several ways to watch a Flink job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks.

FAQ

What is the default target file size in Iceberg?

536,870,912 bytes, 512 MB, set by the table property write.target-file-size-bytes. The rewrite_data_files procedure uses the same value as its default target, and a Flink write can override it with the target-file-size-bytes write option.

How often should a streaming Iceberg table be compacted?

The Iceberg documentation gives no fixed schedule. It names the files metadata table as the place to decide when to compact a partition, so the usual approach is to compact when a partition’s file count or average file size crosses a chosen threshold. The Flink TableMaintenance API can trigger a rewrite after a set number of new data files. At fleet scale, the AutoComp paper describes an automatic compaction framework that balances the benefits of compaction against its costs, and the LinkedIn use case covers how it is used.

Does compaction reduce storage straight away?

No. Compaction writes new files and commits a new snapshot, and the old small files stay referenced by earlier snapshots until those snapshots expire. Storage falls after expire_snapshots runs, as long as the last snapshot of a Flink writer is kept.

Why not just set a very long checkpoint interval?

A longer interval does give larger files, but data reaches the table less often and the Flink job restores from an older checkpoint after a failure. Distribution mode, writer parallelism and compaction reduce small files without lengthening the commit interval, and each carries its own cost, listed in the sections above.

Is Iceberg compaction the same as Kafka log compaction?

No. Kafka log compaction is run by the broker’s log cleaner, which recopies log segments and removes records whose key appears again later in the log, as the Kafka design documentation describes. A record with a null value, called a tombstone, marks earlier records with that key for removal, and the removal happens when the cleaner recopies the segment. Iceberg compaction rewrites a table’s data files into larger ones through a Spark or Flink job.

Related reading