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
- Set the GBK output iterable coder's element coder ID to the input KV's value coder ID
- Rebuild the pipeline from Beam code so value coders are inferred identically on both sides
- Review custom coder rewrite logic to preserve value coder identity through GBK
- 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
- Point the GBK output iterable coder's element at the input KV's value coder ID
- Preserve coder IDs verbatim through graph rewrites
- Unit-test graph-rewrite tools asserting value coder identity across GBK
- Serialize and consume pipelines with one SDK version
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
- Bad coder for input of
- Bad coder for output of
- Input key coder does not match output key coder for…
- Output value coder for transform must be an iterable coder…
- Encountered a type that is not currently supported by…
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)