Skip to content

Issue #8708 : Add Beam MQTT, Debezium, Elasticsearch, Splunk and Snowflake IO - #8716

Open
mattcasters wants to merge 20 commits into
apache:mainfrom
mattcasters:issue-8708
Open

mattcasters wants to merge 20 commits into
apache:mainfrom
mattcasters:issue-8708

Conversation

@mattcasters

Copy link
Copy Markdown
Contributor

This continues the Beam engine work for #8708. The new commit adds the remaining IO transforms and the source-topology guard that goes with them.

Beam MQTT input and output, Debezium input, and Elasticsearch input and output are Beam transforms with annotated dialogs. A source rejects an enabled incoming hop, including an error hop, before it expands. Splunk output writes HEC events through SplunkIO. Failure text contains the HTTP status only. Snowflake input and output copy a bounded collection through a GCS external stage. There is no streaming read and no Snowpipe write. The Direct runner run configuration can set the GCP temp location separately from the general temp location.

The Beam plugin non-UI suite ran 307 tests and the isolated-display UI suite ran 187 tests, both with no failures. The Beam module RAT check reported 0 unapproved files. No pipeline was run against a live Snowflake account or GCS bucket, and the new connectors are not certified on remote runners. Parquet and Avro file handlers are still not implemented. The Flink job-id finding is documented and has not been posted on #3131.

addresses #8708


Thank you for your contribution! Follow this checklist to help us incorporate your contribution quickly and easily:

  • Run mvn clean install apache-rat:check to make sure basic checks pass. A more thorough check will be performed on your pull request automatically.
  • If you have a group of commits related to the same change, please squash your commits into one and force push your branch using git rebase -i.
  • Mention the appropriate issue in your description (for example: addresses #123), if applicable.

To make clear that you license your contribution under the Apache License Version 2.0, January 2004
you have to acknowledge this by using the following check-box.

Track all 18 sub-issues of the Beam engine / runners / Beam IO umbrella
with a decision, effort rating and rationale each, plus the open
questions that block wave 4.

Refs apache#8708
PipelineTestBase runs a pipeline but discards the engine and leaves the
output folder untouched, so tests extending it can neither read the output
nor assert on a row count.  Every Beam transform handler in umbrella apache#8708
needs exactly those two assertions.

Adds SingleTransformPipelineTestBase with helpers to clear stale output,
read the output files, lines, text and per-shard bytes, and read a named
Beam counter off the engine, plus a self-test that exercises all of them.

Two behaviours the helpers have to respect, both found by the self-test:

- Beam shards file output, so a batch arrives as several
  name-00000-of-000NN files.  A helper assuming a single file silently
  drops rows.
- Which counter a transform reports depends on its role: the file input
  counts 'input' and 'written', the file output counts 'read' and
  'output'.  The metric name is a parameter rather than hard coded.

The engine publishes metrics from a background timer, so a read straight
after waitUntilFinished can come back empty; the counter helpers poll
instead of racing, and the negative tests use a non-polling read so they
fail fast.

Refs apache#8708
The type derivation did AvroType.valueOf("RECORD"), which threw and
failed the entire read for any table with a nested or repeated field.
Map RECORD and STRUCT to a String and render the value as JSON.

A field the table schema does not mention (common for a fromQuery
result) also threw "Unable to find field"; it now falls back to the Hop
type declared in the transform and logs a warning.

Nested values are rendered with GenericData.toString(), not Jackson:
Jackson serialises a GenericRecord as its internal bean structure
({"schema":...,"elementType":...}) rather than its contents, which is
what an earlier attempt produced.  A test guards against that.

Fixes apache#5064
…tions

BigQueryIO.TypedRead already offers withQueryLocation() and
withoutValidation(); neither was reachable from Hop.  Both are exposed on
the BigQuery input transform, passed through to the read, and default to
off so a transform saved before this change behaves exactly as before.

Applies only when a query is set, since a plain table read has neither a
query location nor a query to validate.

Fixes apache#2416
…ansform

TextIO.Write.withCompression() already supports it; nothing in Hop
exposed it.  Adds a compression option to the transform, the dialog and
the PTransform, defaulting to blank so existing transforms are unchanged.

Beam's AUTO is read-only, it fails a write with "AUTO is not supported for
writing", so AUTO is resolved against the file suffix here instead of
being passed through.  With no recognisable suffix it degrades to
uncompressed rather than failing the job.

Covered end to end on the Direct runner: the tests run a real pipeline
and check the gzip magic on disk, then gunzip it and check all 100 rows
survive with all 10 fields.

Fixes apache#2337
…onfiguration

BigQuery load jobs and the other GCP IOs need a Cloud Storage path in
GcpOptions.gcpTempLocation, which is a separate Beam option from the
general tempLocation.  GcpOptions is already on the plugin classpath, so
this is a new run configuration option plus the wiring.

Only set when non-blank, so an existing run configuration keeps exactly
the behaviour it had.

The test builds PipelineOptions directly rather than going through
getPipelineOptions(), because that call asks the Dataflow client for GCP
application default credentials, which a unit test cannot have.

Fixes apache#2355
Packaging: the Confluent serializers were only on the compile classpath,
so the Hop fat jar never carried them.  FatJarBuilder stages from
plugins/engines/beam/lib, which did not exist at all, and a Flink task
manager failed with ClassNotFoundException:
io.confluent.kafka.serializers.AbstractKafkaSchemaSerDe.  A
copy-dependencies execution now stages the io.confluent jars there at
prepare-package.

No new dependency declarations: the jars already resolve transitively at
the version the hop-libs BOM manages, and declaring them explicitly
breaks resolution because they are not on Maven Central and the project
only configures central.

Producer: an Avro message field with no schema made the transform call
AvroCoder.of(null), which throws a bare NullPointerException from inside
Beam's determinism checker before the pipeline is submitted.  The check
now names the field and says what to do about it.

The staging test is verified to fail without the copy-dependencies
execution.  The generated jars are covered by the root *.jar gitignore.

Fixes apache#2675
Records the commit for each of the five issues implemented so far.
The 12 remaining rows stay Planned and apache#2613 stays Blocked on scope.

Refs apache#8708
Without a handler the transform fell through to BeamGenericTransformHandler,
which runs the Hop transform inside a ParDo per worker, so every worker
generated its own sequence starting at the same number.

The handler funnels all rows onto one key and keeps the counter in Beam
state, which is what makes it a single pipeline-wide sequence. A plain
field would restart per worker, and state is per key, so the keying and
the state belong together.

Three things worth remembering for any future stateful Beam DoFn here:

- @stateid is a field-only annotation and belongs on a final StateSpec
  field. The State itself arrives as a @ProcessElement parameter, not as
  a field; a field for it is never injected and stays null.
- A stateful ParDo needs a KvCoder input, so the rows are re-keyed after
  Flatten, which leaves a PCollection<Iterable<HopRow>>.
- The row meta the converter passes already carries the sequence field,
  because AddSequenceMeta.getFields() added it. Adding another one grows
  the row past what the output transform reads and silently drops the
  value, which is what an earlier version did.

The database sequence mode cannot be expressed on Beam and now fails with
an explanation rather than emitting duplicates. A trade-off worth
stating: the transform runs at parallelism one, which is what a globally
increasing sequence costs.

The end-to-end test is verified to fail without the handler registered.

Fixes apache#2379
Beam can only group keyed elements, so windowing per key means keying the
rows, windowing the KV collection, grouping, and dropping the key again.
The key is derived from a field that is already on the row, so no extra
column appears in the output.

An empty key field keeps the previous global windowing, which is what a
transform saved before this change gets.

Two details that cost time and are easy to get wrong again:

- WithKeys.of is overloaded: one takes a constant key, the other a
  SerializableFunction. HopKeyFn is the latter, and the call site names the
  type explicitly so the function overload is selected.
- A Window<HopRow> cannot be applied to PCollection<KV<String, HopRow>>.
  Window is a PTransform, not a WindowFn, so the window function is taken
  from window.getWindowFn() and re-wrapped for the keyed element type.

A missing key field now fails the pipeline naming the field, rather than
windowing everything together and looking like it worked. Note that a
valid key field and no key field produce the same row count for a GLOBAL
window, so the end-to-end test proves the keyed branch is reached via that
failure case rather than by row count alone.

The Beam window dialog gets a Key field text box; TextVar has no
setToolTip, the SWT method is setToolTipText.

Fixes apache#2275
Wraps Beam's Partition transform, giving Beam pipelines the partitioning
option the local engine has. This also gives the existing SinglePartitionFn
a caller: it was on the classpath with nothing referencing it.

Two modes: Single sends everything to one partition, Key hashes a field so
rows sharing a value land together.

The part worth reading twice: Beam returns a PCollectionList, one collection
per partition, and the partition for a given row is only decided while the
pipeline runs, so the next transform cannot branch on it. Reading partition
0 therefore drops every row that went to the others - 88 of 100 in the test.
The partitions are flattened back into one stream instead. The end-to-end
test asserts the row count and was verified to fail (12 of 100) against the
single-partition version, so this regression cannot come back quietly.

KeyedPartitionFn carries the value meta as JSON rather than as an IValueMeta,
because Beam serialises the partition function to the workers and an
IValueMeta does not survive that. The meta is rebuilt on first use. A null
key and a zero partition count are both guarded rather than left to throw
inside the runner, and the modulo uses floorMod because Math.abs returns a
negative number for Integer.MIN_VALUE, which would give an invalid index.

Fixes apache#2040
… the

transforms need work

The issue asks whether adding a handler is worth the effort. Investigated it
rather than guessing, and the answer is that the effort splits: two of the
four Avro transforms need nothing, and two need real work. That is a larger
change than the original triage assumed, which had called it handler code
only.

The transport is already done. HopRowCoder handles GenericRecord explicitly,
writing the schema as JSON next to the binary payload - which is exactly
what Avro input needs, since the schema is only known once the file has been
read. Someone anticipated this. HopRowCoderAvroTest now covers the round trip
including a nested record, and pins a detail that is easy to trip over: the
coder is type-tagged, not Java serialization, so a row must hold Long rather
than Integer or it is rejected before the Avro field is reached.

AvroEncode and AvroDecode are pure per-row conversions with no cross-worker
state, which is what the generic handler provides. They work as they are.

AvroFileInput and AvroOutput cannot work there. The input is not a source: it
reads the filename to open from a field in an incoming row, so on Beam there
is no input collection, the generic handler injects one dummy HopRow, and the
transform fails asking for a filename field. The output opens a single writer
and holds it for every row, so every worker would write to the same file.

Both need a handler built on AvroIO, which is already on the classpath. One
complication for whoever picks this up: hop-tech-avro is not a dependency of
the Beam module, so the meta classes cannot be referenced at compile time
without adding it.

Recorded in the triage table. No handler written yet.

Answers apache#2358
…andler needs

The issue asks whether a handler is worth the effort, so this investigates
rather than implements.

The missing dependency was real: ParquetIO ships in
beam-sdks-java-io-parquet, not beam-sdks-java-core, and the build fails to
compile ParquetIOViabilityTest without it.

The triage expected the only remaining work to be handler code. That is not
quite right, though the mismatch is in our favour. ParquetIO is Avro-typed:
read(schema) and sink(schema) both take an org.apache.avro.Schema rather
than a Parquet MessageType, and yield GenericRecord. That looks like an
obstacle but is actually a clean seam, because ParquetOutput.buildSchema()
already builds SchemaBuilder.record("ApacheHopParquetSchema") and derives
the Parquet type from it. Both sides already agree on the schema, so a
handler has no translation layer to invent.

Verified along the way:

- AvroCoder is in org.apache.beam.sdk.extensions.avro.coders, not
  org.apache.beam.sdk.coders.
- GenericRecord.equals is reference-based, so AvroCoder reports
  consistentWithEquals() false even for value-equal records. Irrelevant here
  and not a handler decision, since ParquetIO sets the coder itself.
- Packaging needs no change. The Parquet jars do not appear in the beam
  plugin's lib/ because that directory is scoped to io.confluent by the
  apache#2675 staging execution, but Beam SDKs are unpacked into lib/core, which
  is already on the classpath; lib-beam is legacy and no longer produced.

The transforms still need a handler: ParquetFileInput reads its filename
from an incoming row field and ParquetFileOutput holds a single writer, so
neither works through the generic handler. And hop-tech-parquet is not a
Beam module dependency, so the meta classes cannot be referenced at compile
time - the same constraint as apache#2358.

151 tests pass. Recorded in the triage table; no handler written yet.

Answers apache#2357
…d Snowflake IO

Sources reject an incoming hop before they expand. Snowflake copies through a
GCS stage and does not implement Snowpipe or a streaming read. The Direct
runner run configuration can set the GCP temp location.
- Bind stage-confluent-serializers execution to generate-resources and clean lib/ directory on mvn clean
- Exclude redundant bcprov-jdk15on dependency in Snowflake IO
- Use RowDataUtil.resizeArray and indexed field assignment in AddSequenceFn
- Improve KeyedPartitionFn hashing for byte arrays and BigDecimals with scale normalization and IValueMeta fallback
- Refactor BeamPartitionDialog to use GuiCompositeWidgets.addScrolledComposite and @GuiWidgetElement annotations
- Remove duplicate license header in BeamDirectPipelineRunConfiguration
- Add unit tests for KeyedPartitionFn hashing and BeamPartitionDialog UI interactions
apache#8708)

- In GuiRegistry, use guiPluginClassName directly instead of calling getClass().getName() on it, preventing filter IDs from erroneously starting with 'java.lang.String'.
- In BaseGuiContextHandler, resolve the target filter object from the context handler instance (via getter, field, or self-reference) before falling back to getInstance() and reflection construction. This fixes NullPointerException when evaluating action filters in scenes/tests or when getInstance() returns null.
- Add unit tests for filter ID registration in GuiRegistryTest and resolution logic in BaseGuiContextHandlerTest.
@mattcasters mattcasters added this to the 2.21 milestone Oct 1, 2026
@mattcasters
mattcasters marked this pull request as ready for review October 1, 2026 12:48
@bamaer

bamaer commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

Blockers

  • Add Sequence fails on a real Hop install. The new BeamAddSequenceTransformHandler creates AddSequenceMeta, but plugins/engines/beam/src/main/resources/dependencies.xml has no ../../transforms/addsequence entry. Every Beam pipeline with Add Sequence fails with NoClassDefFoundError: AddSequenceMeta, including the existing beam_directrunner/main-0003 integration test. Unit tests and CI don't catch it: Maven puts all plugins on one classpath, and CI doesn't run the integration tests. Fix: add <folder>../../transforms/addsequence</folder>. With that line, all 15 beam_directrunner integration tests pass.

Should fix

  • Add Sequence on Beam regresses three cases that worked in 2.18. The new handler does fix duplicate values in batch pipelines: 2.18 gave 582 distinct values for 1000 rows, this PR gives 1000. But it breaks three cases that worked in 2.18:
    • Unbounded input (Kafka, MQTT, Debezium, a never-ending Row Generator) now fails with GroupByKey cannot be applied to non-bounded PCollection in the GlobalWindow without a trigger.
    • "Use DB to get sequence" is refused (BeamAddSequenceTransformHandler.java:103).
    • Rows above the max value are dropped without an error (AddSequenceFn.java:125). With max 10 and 30 rows in, 10 rows come out; 2.18 and the local engine wrap around.
    • Fix: use the generic handler for DB mode and unbounded input, and wrap around like Counter.
  • The keyed window ignores the trigger, allowed lateness and discarding mode (BeamWindowMeta.java:313). The keyed path rebuilds the window from getWindowFn() only. A Global window with a trigger on an unbounded input works without a key, but fails with the GroupByKey error above once a key field is set. Fix: apply the same settings to keyedWindow.
  • Compressed file output doubles the extension. With suffix .csv.gz, GZIP and AUTO both write rows-00000-of-00003.csv.gz.gz (BeamOutputTransform.java:116,139). The docs tell users to put .gz in the suffix for AUTO, so AUTO always doubles it. Fix: strip the codec's suffix before withSuffix, and assert file names in the test.
  • The new dependencies add 114 MB to lib/core for everyone who installs Beam.
    • 96.5 MB of that is snowflake-jdbc-4.3.0. It contains its own unshaded copy of zstd-jni, and with Beam installed Hop's Parquet ZSTD codec loads com.github.luben.zstd.Zstd from the Snowflake jar instead of zstd-jni-1.5.7-13.jar. It works, but which copy wins depends on classpath order.
    • beam-sdks-java-io-parquet (pom.xml:1176) is compile scope, but only ParquetIOViabilityTest uses it, so about 10 MB of Parquet jars ship unused. Make it test scope.
    • Consider whether SnowflakeIO belongs in the default Beam plugin.
  • The Beam partition transform has no effect. BeamPartitionMeta.java:215 flattens the partitions back into one stream, so a following GroupByKey still shuffles. beampartition.adoc claims it avoids that shuffle and that "single partition serialises processing". Fix the docs, or give each partition its own output.
  • LICENSE (line 520) and NOTICE (line 52) list bcprov-jdk15on, but the pom excludes it and it's not in the zip. Remove those entries; the NOTICE paragraph about licenses belongs in LICENSE.
  • issue-8708-triage.adoc is project tracking in the user manual, with commit hashes that won't survive a squash, and it's linked from nav.adoc:70. Move it to the umbrella issue.

Verified: built with all unit tests; ran the beam_directrunner integration tests; ran pipelines on the Direct runner against live MQTT, Elasticsearch 8, Splunk 9 HEC and Postgres/Debezium; compared Add Sequence behaviour with Hop 2.18. Not run: Snowflake, BigQuery, remote runners.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants