apache/beam · error · ValueError

Input value coder does not match output value coder for…

Error message

Input value coder %s does not match output value coder %s for transform %s

What it means

Raised by validate_transform when the value component coder of a GROUP_BY_KEY's output does not match the value coder of its input. GBK preserves value encodings (output is KV<K, Iterable<V>> over the same V encoding as the input KV<K,V>); a mismatch indicates corrupted or rewritten coder wiring in the proto.

Solutions

  1. Set the GBK output iterable coder's element coder ID to the input KV's value coder ID
  2. Rebuild the pipeline from Beam code so value coders are inferred identically on both sides
  3. Review custom coder rewrite logic to preserve value coder identity through GBK
  4. Re-serialize the pipeline with one consistent SDK version

Example fix

// before
iter_coder = make_iterable_coder(other_value_coder_id)
// after
iter_coder = make_iterable_coder(input_coder.component_coder_ids[1])
Defensive patterns

Strategy: validation

Validate before calling

def check_value_coders_match(pipeline, t) -> bool:
    if t.spec.urn != common_urns.primitives.GROUP_BY_KEY.urn:
        return True
    in_c = get_coder(next(iter(t.inputs.values())))
    out_c = get_coder(next(iter(t.outputs.values())))
    values_coder = pipeline.components.coders[out_c.component_coder_ids[1]]
    return values_coder.component_coder_ids[0] == in_c.component_coder_ids[1]

Type guard

def values_consistent(pipeline, t) -> bool:
    in_c, out_c = gbk_coders(pipeline, t)
    element_id = pipeline.components.coders[out_c.component_coder_ids[1]].component_coder_ids[0]
    return element_id == in_c.component_coder_ids[1]

Try / catch

try:
    validate_pipeline_graph(pipeline_proto)
except ValueError as e:
    if 'does not match output value coder' in str(e):
        point_iterable_at_input_value_coder(pipeline_proto)
    else:
        raise

Prevention

When it happens

Trigger: output_values_coder.component_coder_ids[0] (the iterable's element coder) != input_coder.component_coder_ids[1] for a GBK transform — usually after custom coder ID reassignment or graph merging.

Common situations: Pipeline graph rewrite/optimization tools reassigning coder IDs; manually stitched pipeline fragments; cross-SDK serialization with divergent coder registries.

Understand the failure class

Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.

Related errors


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

Appendix: source

Thrown at sdks/python/apache_beam/runners/pipeline_utils.py:116

      if input_key_coder_id != output_key_coder_id:
        raise ValueError(
            "Input key coder %s does not match output key coder %s for "
            "transform %s" %
            (input_key_coder_id, output_key_coder_id, transform_id))
      output_values_coder_id = output_coder.component_coder_ids[1]
      output_values_coder = pipeline_proto.components.coders[
          output_values_coder_id]
      if output_values_coder.spec.urn != common_urns.coders.ITERABLE.urn:
        raise ValueError(
            "Output value coder %s for transform %s must be an iterable "
            "coder, but uses URN %s" % (
                output_values_coder_id,
                transform_id,
                output_values_coder.spec.urn))
      input_value_coder_id = input_coder.component_coder_ids[1]
      output_value_coder_id = output_values_coder.component_coder_ids[0]
      if output_value_coder_id != input_value_coder_id:
        raise ValueError(
            "Input value coder %s does not match output value coder %s for "
            "transform %s" %
            (input_value_coder_id, output_value_coder_id, transform_id))
    elif transform_proto.spec.urn == common_urns.primitives.ASSIGN_WINDOWS.urn:
      if not transform_proto.inputs:
        raise ValueError("Missing input for transform: %s" % transform_proto)
    elif transform_proto.spec.urn == common_urns.primitives.PAR_DO.urn:
      if not transform_proto.inputs:
        raise ValueError("Missing input for transform: %s" % transform_proto)

    for t in transform_proto.subtransforms:
      validate_transform(t)

  for t in pipeline_proto.root_transform_ids:
    validate_transform(t)


def _dep_key(dep):

View on GitHub (pinned to 12126d8942)