apache/beam · error · RuntimeError

Ambiguity due to non-preserved tags

Error message

Ambiguity due to non-preserved tags: %s vs %s

What it means

When re-wiring an externally expanded transform back into the pipeline graph, an input tag no longer matches after tag renaming and this happened for a second time; the runner cannot tell which PCollection the renamed tag corresponds to, so the mapping is ambiguous.

Solutions

  1. Report to Beam devs: non-preserved tags are an expansion-service protocol edge case
  2. As a workaround, avoid multi-input external transforms with duplicated/renamed tags
Defensive patterns

Strategy: fallback

When it happens

Trigger: Thrown at sdks/python/apache_beam/transforms/external.py:921 when the library encounters an invalid state.

Common situations: See trigger scenarios.


AI-assisted analysis of apache/beam@12126d8942 (2026-09-13). Data as JSON: /api/errors/455f3d5b4ac5ffdf. Report an issue: GitHub.

Appendix: source

Thrown at sdks/python/apache_beam/transforms/external.py:921

      return env

    for env in components.environments.values():
      _resolve_artifacts_for(env)
    return components

  def _output_to_pvalueish(self, output_dict):
    if len(output_dict) == 1:
      return next(iter(output_dict.values()))
    else:
      return output_dict

  def to_runner_api_transform(self, context, full_label):
    pcoll_renames = {}
    renamed_tag_seen = False
    for tag, pcoll in self._inputs.items():
      if tag not in self._expanded_transform.inputs:
        if renamed_tag_seen:
          raise RuntimeError(
              'Ambiguity due to non-preserved tags: %s vs %s' % (
                  sorted(self._expanded_transform.inputs.keys()),
                  sorted(self._inputs.keys())))
        else:
          renamed_tag_seen = True
          tag, = self._expanded_transform.inputs.keys()
      pcoll_renames[self._expanded_transform.inputs[tag]] = (
          context.pcollections.get_id(pcoll))
    for tag, pcoll in self._outputs.items():
      pcoll_renames[self._expanded_transform.outputs[tag]] = (
          context.pcollections.get_id(pcoll))

    def _equivalent(coder1, coder2):
      return coder1 == coder2 or _normalize(coder1) == _normalize(coder2)

    def _normalize(coder_proto):
      normalized = copy.copy(coder_proto)
      normalized.spec.environment_id = ''

View on GitHub (pinned to 12126d8942)