fix(NuEventEmitter): allow Delta MERGE events through the job-name allowlist - #22
Merged
Merged
Conversation
…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>
There was a problem hiding this comment.
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 inNuEventEmitter’sWANTED_EVENT_NAME_SUBSTRINGS. - Add a new
NuEventEmitterTestcovering 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.
lmassaoy
approved these changes
Sep 11, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
WANTED_EVENT_NAME_SUBSTRINGSinNuEventEmitter.javaonly 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 innubank/data-lineage's enhancement app (core/enhancement.py, fixed separately in data-lineage#223).IncrementalMergeJobmerge step is missing from Marquez org-wide, and in turn why merge-fed novelty tables show a nullout_edgein the curated lineage datasets.execute_merge_into_commandvia this agent's ownNameNormalizer) are on MDP-566.data-lineage's copy of this allowlist is fixed in a separate PR (linked above);spark-runtimeneeds itsnu-openlineage-spark-agentpins 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..execute_merge_into_command.— Iceberg's.replace_data./.write_delta.were considered and dropped:job_interface_scala.job.Formathas no Iceberg case,IncrementalMergeJob's merge path is hardcoded toio.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
NuEventEmitterTest.java(no test existed for this class before) covering: Delta MERGECOMPLETEevents are now emitted; all three previously-allowlisted commands still pass (no regression); non-allowlisted job names, non-SQL_JOBjob types, andRUNNINGevents 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 unmodifiedmain. Pre-existing, unrelated to this change../gradlew publishToMavenLocalfromclient/javaandintegration/sql/iface-javafirst (sibling monorepo modules the Spark integration depends on as1.38.0-SNAPSHOT, not resolvable from Maven Central) to build:appat all.🤖 Generated with Claude Code