This is an automated email from the ASF dual-hosted git repository.
corgy-w pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new d49c5a19b5 [Fix][CLI] Include transform blocks in routing validation
and run it unconditionally (#11658)
d49c5a19b5 is described below
commit d49c5a19b56e26650a3463f997e869b66892a810
Author: SEZ <[email protected]>
AuthorDate: Mon Aug 10 21:47:09 2026 +0800
[Fix][CLI] Include transform blocks in routing validation and run it
unconditionally (#11658)
---
seatunnel-cli/seatunnel_cli/agents.py | 38 ++++---
.../tests/test_validate_hocon_placeholders.py | 118 +++++++++++++++++++++
2 files changed, 139 insertions(+), 17 deletions(-)
diff --git a/seatunnel-cli/seatunnel_cli/agents.py
b/seatunnel-cli/seatunnel_cli/agents.py
index 14caee8a04..07adfc6670 100644
--- a/seatunnel-cli/seatunnel_cli/agents.py
+++ b/seatunnel-cli/seatunnel_cli/agents.py
@@ -288,7 +288,7 @@ def _extract_connector_blocks_raw(
(inclusive), e.g. ``'url = "..." driver = "..."'``.
"""
results: list[tuple[str, str, str]] = []
- for section in ("source", "sink"):
+ for section in ("source", "transform", "sink"):
# Find the section's opening brace
pattern = re.compile(
rf"(?:^|\n)\s*{section}\s*\{{", re.IGNORECASE,
@@ -345,8 +345,13 @@ def _validate_routing_pairs(
errors: list[str],
warnings: list[str],
) -> None:
- """Validate plugin_output / plugin_input pairing across connector
blocks."""
- outputs: dict[str, str] = {} # label → connector_name
+ """Validate plugin_output / plugin_input pairing across connector blocks.
+
+ Blocks may come from the source, transform, or sink sections; messages
+ carry the real ``section.connector`` location so diagnostics point at
+ the right block (transforms both emit and consume labels).
+ """
+ outputs: dict[str, str] = {} # label → "section.connector_name"
inputs: dict[str, str] = {}
label_re = re.compile(
@@ -354,18 +359,19 @@ def _validate_routing_pairs(
)
for section, connector_name, block_content in blocks:
+ location = f"{section}.{connector_name}"
for m in label_re.finditer(block_content):
label = m.group(1)
if "plugin_output" in m.group(0):
if label in outputs:
errors.append(
f"Duplicate plugin_output label \"{label}\" "
- f"in source.{connector_name} "
- f"(already used by source.{outputs[label]})"
+ f"in {location} "
+ f"(already used by {outputs[label]})"
)
- outputs[label] = connector_name
+ outputs[label] = location
else:
- inputs[label] = connector_name
+ inputs[label] = location
# Only validate pairing when routing labels are actually used
if not outputs and not inputs:
@@ -375,16 +381,16 @@ def _validate_routing_pairs(
if unmatched_inputs:
for label in unmatched_inputs:
errors.append(
- f"sink.{inputs[label]}: plugin_input \"{label}\" "
- f"has no matching plugin_output in any source"
+ f"{inputs[label]}: plugin_input \"{label}\" "
+ f"has no matching plugin_output in any source or transform"
)
orphan_outputs = set(outputs) - set(inputs)
if orphan_outputs:
for label in orphan_outputs:
warnings.append(
- f"source.{outputs[label]}: plugin_output \"{label}\" "
- f"has no matching plugin_input in any sink"
+ f"{outputs[label]}: plugin_output \"{label}\" "
+ f"has no matching plugin_input in any transform or sink"
)
@@ -463,9 +469,6 @@ def validate_hocon(config_str: str) -> str:
except Exception:
pass
- # Validate routing labels
- _validate_routing_pairs(raw_blocks, errors, warnings)
-
elif parsed is not None:
# Standard single-connector-per-type path (pyhocon is accurate here)
for section in ["source", "sink"]:
@@ -496,9 +499,10 @@ def validate_hocon(config_str: str) -> str:
except Exception:
pass
- # Also validate routing for single-connector configs if labels present
- if raw_blocks:
- _validate_routing_pairs(raw_blocks, errors, warnings)
+ # Routing validation is independent of option metadata and pyhocon —
+ # always run it when connector blocks were extracted.
+ if raw_blocks:
+ _validate_routing_pairs(raw_blocks, errors, warnings)
# Check STREAMING mode needs checkpoint.interval
if parsed is not None:
diff --git a/seatunnel-cli/tests/test_validate_hocon_placeholders.py
b/seatunnel-cli/tests/test_validate_hocon_placeholders.py
index 4e9a32a634..3b14f26513 100644
--- a/seatunnel-cli/tests/test_validate_hocon_placeholders.py
+++ b/seatunnel-cli/tests/test_validate_hocon_placeholders.py
@@ -103,3 +103,121 @@ def
test_colon_separator_still_rejects_env_vars_elsewhere():
config = _config_with_sink('topic: "${now}"')
result = validate_hocon(config)
assert "Unresolved environment variables" in result
+
+
+# ── transform-mediated routing (regression for transform blocks being ──
+# ── invisible to _validate_routing_pairs) ──
+
+def test_sink_consuming_transform_output_is_valid():
+ config = """
+env { parallelism = 1
+ job.mode = "BATCH" }
+source { FakeSource {
+ row.num = 5
+ schema { fields { id = "bigint" } }
+ plugin_output = "raw" } }
+transform {
+ Sql {
+ plugin_input = "raw"
+ plugin_output = "filtered"
+ query = "SELECT * FROM raw WHERE id > 1"
+ }
+}
+sink { Console { plugin_input = "filtered" } }
+"""
+ result = validate_hocon(config)
+ assert "has no matching plugin_output" not in result
+
+
+def test_split_via_parallel_transforms_is_valid():
+ config = """
+env { parallelism = 1
+ job.mode = "BATCH" }
+source { FakeSource {
+ row.num = 5
+ schema { fields { id = "bigint" } }
+ plugin_output = "raw" } }
+transform {
+ Sql {
+ plugin_input = "raw"
+ plugin_output = "big"
+ query = "SELECT * FROM raw WHERE id > 100"
+ }
+ Sql {
+ plugin_input = "raw"
+ plugin_output = "small"
+ query = "SELECT * FROM raw WHERE id <= 100"
+ }
+}
+sink {
+ Console { plugin_input = "big" }
+ Console { plugin_input = "small" }
+}
+"""
+ result = validate_hocon(config)
+ assert "has no matching plugin_output" not in result
+
+
+def test_genuinely_unmatched_plugin_input_still_rejected():
+ config = """
+env { parallelism = 1
+ job.mode = "BATCH" }
+source { FakeSource {
+ row.num = 5
+ schema { fields { id = "bigint" } }
+ plugin_output = "raw" } }
+sink { Console { plugin_input = "nonexistent_label" } }
+"""
+ result = validate_hocon(config)
+ assert "nonexistent_label" in result
+ assert "has no matching plugin_output" in result
+
+
+def test_routing_errors_attribute_transform_blocks_correctly():
+ # A dangling plugin_input on a TRANSFORM must be reported against the
+ # transform block, not misattributed to a source or sink.
+ config = """
+env { parallelism = 1
+ job.mode = "BATCH" }
+source { FakeSource {
+ row.num = 5
+ schema { fields { id = "bigint" } }
+ plugin_output = "raw" } }
+transform {
+ Sql {
+ plugin_input = "wrong_label"
+ plugin_output = "filtered"
+ query = "SELECT * FROM wrong_label"
+ }
+}
+sink { Console { plugin_input = "filtered" } }
+"""
+ result = validate_hocon(config)
+ assert 'transform.Sql: plugin_input "wrong_label"' in result
+
+
+def test_duplicate_output_across_transforms_reports_transform_location():
+ config = """
+env { parallelism = 1
+ job.mode = "BATCH" }
+source { FakeSource {
+ row.num = 5
+ schema { fields { id = "bigint" } }
+ plugin_output = "raw" } }
+transform {
+ Sql {
+ plugin_input = "raw"
+ plugin_output = "same_label"
+ query = "SELECT * FROM raw WHERE id > 1"
+ }
+ Sql {
+ plugin_input = "raw"
+ plugin_output = "same_label"
+ query = "SELECT * FROM raw WHERE id <= 1"
+ }
+}
+sink { Console { plugin_input = "same_label" } }
+"""
+ result = validate_hocon(config)
+ assert 'in transform.Sql' in result
+ assert 'already used by transform.Sql' in result