Hi Rui, hi Richard,

Both of you are right, and on the two points that matter I'm changing the
proposal rather than defending it.

On Trident: the objection lands harder than intended, because my own
opening mail argued for tuple-at-a-time control over what is durably
persisted and then proposed a micro-batch API. So it becomes a bolt.
Catalog config, writer factory, partitioned fanout and commit logic have no
Trident types in them anyway.

On the WAL: I don't see a hole, and working through it separated two things
I had conflated. The WAL gives commit atomicity, since the data files are
durable before the commit and only the reference to them needs protecting.
What Trident was supplying was something different, a deterministic
identity of the input, and that comes from the source, not from the sink. A
general-purpose module can't assume users have a replayable,
deterministically addressed source, so exactly-once isn't mine to promise.

What I'll claim instead, and document in these words: atomic commits with
at-least-once delivery. Readers never see a partial batch, a crash costs
orphan data files rather than a broken table, and duplicates are possible
on replay with no equality deletes to remove them, so they stay visible
until something downstream dedupes. An extractor SPI for sources that can
do better is a v2 conversation, not a v1 claim.

Everything else I simply agree with: properties-only per batch is dropped,
it buys nothing; not bundled and no transitive catalog implementations, per
#8819; explicit Iceberg pin against the Java baseline; format v2 and
append-only stated plainly; maintenance out of scope but documented as
still necessary. On ownership, yes, I'm signing up to keep it current, and
I'll say so openly if that stops being true rather than let it rot.

Prototype of the bolt and the WAL path next, then moving PR from draft to
ready.

Thanks, this is a better proposal than the one I started with,
Gianluca

Il giorno mar 28 lug 2026 alle ore 19:28 Richard Zowalla <[email protected]>
ha scritto:

> Hi Gianluca,
>
> I am not a Trident nor Iceberg user, so take the below with that caveat
> because it's mostly design and project-level questions, not "I've run this
> in production" feedback.
>
> On the HdfsState analogy: Worth being precise about what HdfsState
> actually does, because it's not quite a txid marker on the target. It
> writes a separate, per-partition index file (.index.<topology>.<partition>)
> next to the data, updated during beginCommit, holding the last-seen txid
> plus a file offset; on replay it truncates back to that offset. Two things
> follow for storm-iceberg:
>
> - A single table property is table-global, but Trident state is
> per-partition. If more than one Trident partition writes the same Iceberg
> table, one storm.txid property can't tell you which partitions have
> committed. You'd want per-partition keys
> (storm.txid.<topology>.<partition>), and it'd be good to say so explicitly
> in the PR.
> - N partitions each doing a catalog commit per batch means N concurrent
> optimistic-concurrency commits against the same table metadata. Appends
> retry cheaply, but it's still N catalog round-trips per batch. This is the
> thing that makes commit batching not really optional at any real rate.
>
> On properties-only-per-batch: I agree with the objection already raised:
> it doesn't help. The data files are written per batch either way, so the
> small-file problem is unchanged, and you still pay the catalog round-trip,
> which is the expensive half. It buys nothing over what you have.
>
> On the buffering trade-off: I don't think the loss window is inherent to
> buffering; it's inherent to acking a batch before anything durable records
> that its data files exist. But the data files are already durable in object
> storage before the Iceberg commit; what's lost on a crash is only the
> reference to them. So there's a middle option that mirrors the HdfsState
> precedent more closely than the table property does: before returning from
> commit(), write a small side file (a WAL, effectively) into the table
> location listing txid → pending data fileso snapshot, one cheap PUT. On
> restart, read it and commit anything notyet in the table; dedupe by txid as
> you already do. That keeps exactly-once and lets files accumulate to target
> size. It's more moving parts, so it may well be a v2 thing — but if it's
> viable, I'd rather see the batching path documented as "exactly-once,
> recovery via WAL" than as "opt-in, weaker". Curious whether you see a hole
> in that?
>
> Failing that shipping batching off-by-default and clearly documented is a
> reasonable first iteration, and I wouldn't block on it.
>
> Project-level, and honestly my bigger concern: why Trident-only? Trident
> is lightly used and effectively in maintenance mode; a plain bolt would
> reach far more people. Even if the first cut is Trident, please keep the
> catalog config, writer factory and partitioned-fanout logic free of Trident
> types so a bolt can be layered on later without a rewrite.
>
> Also:
>
> - We recently unbundled optional Hadoop/Kafka deps  binary distribution
> (#8819). Iceberg drags in a large tree. storm-iceberg should not be bundled
> by default, and catalog implementa must be user-supplied rather than
> transitive.
> - Please pin the Iceberg version explicitly and confirm it against our
> Java baseline, and state which table format version(s) you target.
> Append-only means no equality deletes - fine, but say it.
> - Table maintenance (compaction, expire-snapshots) is out of scope and
> should stay out of scope, but the docs should tell users they still need to
> run it.
> - Maintenance commitment: Iceberg moves quickly and we have little Iceberg
> expertise here. Are you signing up to keep this current across Iceberg
> releases? That matters more to me than any of the technical points above -
> we have external modules that nobody owns, and I'd rather not add another
> (we already dropped quite a few in the last months)
>
> None of this is a -1. I'd like to see it land, with the scope questions
>
> Gruß
> Richard
>
> > Am 26.07.2026 um 12:45 schrieb Gianluca Graziadei <
> [email protected]>:
> >
> > Hi folks,
> > I'd like to propose a new external module, *storm-iceberg*, and get the
> > community's thoughts before/while pushing it as a PR.
> >
> > Today the common way to land Storm data into a table format like Iceberg
> is
> > indirect: Kafka plus a separate ingestion job (Spark Structured
> Streaming,
> > Kafka Connect, etc.). That's a fine default for most
> streaming-to-lakehouse
> > pipelines, but it adds a hop and gives up something Storm is unusually
> good
> > at: tuple-at-a-time processing with per-record acking and backpressure.
> For
> > use cases that need real, sub-second control over what has and hasn't
> been
> > durably persisted, not "eventually, in the next micro-batch", Storm can
> > write straight to Iceberg, with nothing in between.
> >
> > What
> > - A Trident State (*IcebergState / IcebergStateFactory /
> > IcebergStateUpdater*) that appends TridentTuples to an Iceberg table.
> > - Exactly-once via the same txid-marker pattern HdfsState already uses:
> > each commit is a single Iceberg transaction that atomically appends the
> > batch's data files and records the transaction id as a table property. A
> > replayed batch is detected and skipped.
> > - Works with any Iceberg catalog (Hadoop, Hive, REST, Glue, Nessie, ...),
> > catalog properties are passed verbatim to
> CatalogUtil.buildIcebergCatalog.
> > - Partitioned tables via Iceberg's PartitionedFanoutWriter.
> > - Optional commit batching: by default every Trident batch is one Iceberg
> > commit and one snapshot, which doesn't scale for high-rate, small-batch
> > streams, thousands of snapshots, one tiny file per batch.
> > *withCommitIntervalBytes
> > / withCommitIntervalMillis* buffer several batches and commit them
> together
> > once a size or time threshold is crossed, so files grow towards the
> table's
> > target size and the snapshot count stays sane.
> >
> > Buffering weakens the delivery guarantee, and I don't think it can be
> made
> > not to. Trident considers a batch delivered the moment commit() returns,
> so
> > a batch buffered towards the next threshold is already acknowledged, if
> the
> > worker dies before the flush, those batches are lost, not duplicated.
> It's
> > opt-in and off by default, documented as such, but I'd welcome opinions
> on
> > whether that's the right trade-off to offer at all, or whether there's a
> > design that keeps exactly-once while still batching commits. (I looked
> at a
> > properties-only commit per batch plus a single deferred append, keeps
> > exactly-once, still costs one catalog round-trip per batch. Happy to
> > discuss both.)
> >
> > Draft PR: https://github.com/apache/storm/pull/8950
> >
> > Feedback, questions, and contributors are all very welcome.
> >
> > Cheers,
> > Gianluca
>
>

Reply via email to