Skip to content

fix(NuEventEmitter): allow Delta MERGE events through the job-name allowlist - #22

Merged
jrosend merged 1 commit into
mainfrom
MDP-566-merge-lineage-allowlist
Sep 11, 2026
Merged

jrosend merged 1 commit into
mainfrom
MDP-566-merge-lineage-allowlist

Conversation

@jrosend

@jrosend jrosend commented Sep 10, 2026

Copy link
Copy Markdown

Summary

  • WANTED_EVENT_NAME_SUBSTRINGS in NuEventEmitter.java only allowed .execute_insert_into_hadoop_fs_relation_command., .adaptive_spark_plan., and .execute_save_into_data_source_command., so every Delta MERGE (.execute_merge_into_command.) event was silently discarded here — independent of the identical allowlist in nubank/data-lineage's enhancement app (core/enhancement.py, fixed separately in data-lineage#223).
  • This is why every IncrementalMergeJob merge step is missing from Marquez org-wide, and in turn why merge-fed novelty tables show a null out_edge in the curated lineage datasets.
  • Full root cause, cross-repo evidence, and a local reproduction (Spark 3.5.3 / Delta 3.3.2 / OpenLineage 1.38.0, matching prod — confirmed a real Delta MERGE produces exactly execute_merge_into_command via this agent's own NameNormalizer) are on MDP-566.
  • Step 2 of 3: data-lineage's copy of this allowlist is fixed in a separate PR (linked above); spark-runtime needs its nu-openlineage-spark-agent pins bumped afterward (both the Scala 2.12/Spark 3.5.3 and Scala 2.13/Spark 4.2.0 pins) once this is published. Shipping this alone does not yet restore merge lineage — see the ticket's validation plan for the full rollout sequence.
  • Scope is intentionally just .execute_merge_into_command. — Iceberg's .replace_data./.write_delta. were considered and dropped: job_interface_scala.job.Format has no Iceberg case, IncrementalMergeJob's merge path is hardcoded to io.delta.tables.DeltaTable, and Iceberg at Nubank is only used via Avalanche's Flink sink connectors, a different stack that doesn't go through this listener.

Test plan

  • Added NuEventEmitterTest.java (no test existed for this class before) covering: Delta MERGE COMPLETE events are now emitted; all three previously-allowlisted commands still pass (no regression); non-allowlisted job names, non-SQL_JOB job types, and RUNNING events are still correctly discarded.
  • ./gradlew :app:test --tests NuEventEmitterTest: 8/8 pass.
  • ./gradlew :app:test (full unit suite, excludes integration/delta/iceberg tags by default): 170 run, 20 failed — confirmed identical 20 failures (same 5 test classes) with this diff stashed out against unmodified main. Pre-existing, unrelated to this change.
  • Required ./gradlew publishToMavenLocal from client/java and integration/sql/iface-java first (sibling monorepo modules the Spark integration depends on as 1.38.0-SNAPSHOT, not resolvable from Maven Central) to build :app at all.

🤖 Generated with Claude Code

…lowlist

WANTED_EVENT_NAME_SUBSTRINGS only allowed
.execute_insert_into_hadoop_fs_relation_command., .adaptive_spark_plan., and
.execute_save_into_data_source_command., so every Delta MERGE
(.execute_merge_into_command.) event was silently discarded here regardless
of jobType/eventType, independent of the identical allowlist in
nubank/data-lineage's enhancement app.

This is why every IncrementalMergeJob's merge step is missing from Marquez
org-wide. Locally reproduced (Spark 3.5.3 / Delta 3.3.2 / this agent's
NameNormalizer) that a real Delta MERGE produces exactly this job-name
segment. See MDP-566.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Copilot AI lite review requested due to automatic review settings September 10, 2026 20:25

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR updates the Spark agent’s Nubank-specific event allowlist so Delta Lake MERGE operations (normalized as .execute_merge_into_command.) are no longer silently filtered out, restoring lineage emission for MERGE steps.

Changes:

  • Allow .execute_merge_into_command. job-name segments in NuEventEmitter’s WANTED_EVENT_NAME_SUBSTRINGS.
  • Add a new NuEventEmitterTest covering MERGE emission and ensuring existing allowlisted commands and filtering rules don’t regress.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated no comments.

File Description
integration/spark/app/src/main/java/io/openlineage/spark/agent/NuEventEmitter.java Extends the job-name substring allowlist to include Delta MERGE events.
integration/spark/app/src/test/java/io/openlineage/spark/agent/NuEventEmitterTest.java Adds unit tests validating MERGE COMPLETE emission and verifying discard behavior for non-allowlisted names, non-SQL job types, and RUNNING events.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

@jrosend
jrosend merged commit 5e7f7cd into main Sep 11, 2026
8 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants