sofiathoren opened a new issue, #40374:
URL: https://github.com/apache/beam/issues/40374
### What happened?
In Apache Beam Python (including `2.76.0`),
`ParDo._get_key_and_window_coder` in `apache_beam/transforms/core.py:1865`
attempts to determine the main input PCollection by running:
```python
main_input = list(set(named_inputs.keys()) - set(self.side_inputs))[0]
```
- `named_inputs.keys()` contains strings (`'None'`, `'python_side_input0'`).
- `self.side_inputs` contains `_SideInput` objects.
Because strings and `_SideInput` objects never match, subtracting
`set(self.side_inputs)` is a no-op: both `'None'` and `'python_side_input0'`
remain in the set.
Due to Python's hash randomization (`PYTHONHASHSEED`), `list(...)[0]`
arbitrarily selects `'python_side_input0'` instead of `'None'`.
When `'python_side_input0'` is selected:
1. `input_pcoll = named_inputs['python_side_input0']`.
2. Because side inputs typically have `element_type` of `Any`,
`kv_type_hint` is `Any`.
3. Beam falls back to `FastPrimitivesCoder` for the timer family key coder
instead of the main input KV key coder.
4. When deployed to a runner such as Dataflow Runner v2 (Windmill), the
runner delivers timer keys encoded with the main input's KV coder (`b"2..."`).
5. The worker bundle processor uses `FastPrimitivesCoderImpl` to decode the
key, treats the first ASCII byte (`0x32` / `'2'`) as an unknown primitive type
tag, and crashes:
```
ValueError: Unknown type tag 32
```
### Steps to reproduce
Create a pipeline with a stateful `DoFn` that also takes a side input:
```python
import apache_beam as beam
from apache_beam.transforms.userstate import ReadModifyWriteStateSpec,
TimerSpec, on_timer
from apache_beam.utils.timestamp import Timestamp
from apache_beam.portability.api import beam_runner_api_pb2
STATE = ReadModifyWriteStateSpec('state', beam.coders.VarIntCoder())
TIMER = TimerSpec('timer', beam.TimeDomain.WATERMARK)
class StatefulFn(beam.DoFn):
def process(
self,
element: tuple[str, int],
state=beam.DoFn.StateParam(STATE),
timer=beam.DoFn.TimerParam(TIMER),
side: dict = None,
):
timer.set(Timestamp.now() + 60)
yield element
@on_timer(TIMER)
def on_timer(self):
yield 1
p = beam.Pipeline()
main_in = p | "Main" >> beam.Create([("key1", 1)])
side_in = p | "Side" >> beam.Create([("s", 1)])
side_dict = beam.pvalue.AsDict(side_in)
main_in | beam.ParDo(StatefulFn(), side=side_dict)
proto = p.to_runner_api()
# Depending on PYTHONHASHSEED, timer family specs will intermittently use
FastPrimitivesCoder for key_coder instead of StrUtf8Coder / KVCoder
```
### Proposed fix
Filter out side input PCollections by object identity, and filter out side
input prefix keys:
```python
def _get_key_and_window_coder(self, named_inputs):
if named_inputs is None or not self._signature.is_stateful_dofn():
return None, None
side_input_pcolls = {si.pvalue for si in (self.side_inputs or ())}
main_input_candidates = [
k
for k, pcoll in named_inputs.items()
if pcoll not in side_input_pcolls
and not str(k).startswith(core.SIDE_INPUT_PREFIX)
]
if main_input_candidates:
main_input = main_input_candidates[0]
else:
main_input = list(named_inputs.keys())[0]
input_pcoll = named_inputs[main_input]
kv_type_hint = input_pcoll.element_type
if kv_type_hint and kv_type_hint != typehints.Any:
coder = coders.registry.get_coder(kv_type_hint)
if not coder.is_kv_coder():
raise ValueError(
f"Input elements to the transform {self} with stateful DoFn
must be "
"key-value pairs."
)
key_coder = coder.key_coder()
else:
key_coder = coders.registry.get_coder(typehints.Any)
window_coder = input_pcoll.windowing.windowfn.get_window_coder()
return key_coder, window_coder
```
### Issue Priority
Priority: 2 (default / most bugs should be filed as P2)
### Issue Components
- Component: Python SDK
- Component: Google Cloud Dataflow Runner
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]