Thanks for sharing. 

Ganesh. 








 ---- On Wed, 23 Sep 2026 05:21:35 +0530  [email protected]  wrote ----

Great question! That's the trick: The Runner isn't expected to handle raw 
pickled KV bytes. It's expected for the runner to tell the SDK worker to use 
different coders when emitting to the runner, (via the data sink), and what the 
SDK worker should expect to receive from the runner via the data source.

That is: The runner is in charge, not the SDK. Portable runners tell the SDK 
what coder they should be using instead in execution. Typically by wrapping the 
coder in question with Length Prefixes, so that the encoded values remain 
distinguishable by the runner.

For KVs and thus GBKs, the key portion is being able to extract the keys for 
the grouping operation.


Typically, this means substituting the existing coder KV<K,V> with the unknown 
components length prefixed.  Eg.if both K and V are unknown coders, then the 
wrapping produces KV<LP<K>, LP<V>>.

You can see how the Prism runner does this here with the lpUnknownCoders 
function: 
https://github.com/apache/beam/blob/master/sdks/go/pkg/beam/runners/prism/internal/coders.go#L134
 



Cheers,
Robert Burke


On Sun, 20 Sept 2026 at 06:41, Ganesh Sivakumar < mailto:[email protected] > 
wrote:

Hi Beam dev community,



I'm currently adding Beam python SDK support for my rust based portable beam 
runner, initially it was designed around to execute java pipelines but, it 
should be able to support majority of python pipeline proto translation logic, 
because of portability. 



However, I've run into a challenge, many Python pipeline pcollections are 
encoded using the pickled coder (beam:coder:pickled_python:v1). I understand 
this is a SDK specific (non-portable) coder and the runner cannot deserialize 
it directly.



how should a runner handle a KV/ tuple encoded this way when performing GBK

For example, KV coder:



ref_Coder_TupleCoder_10": Coder {

                    spec: Some(

                        FunctionSpec {

                            urn: "beam:coder:kv:v1",

                            payload: [],

                        },

                    ),

                    component_coder_ids: [

                        "ref_Coder_DeterministicFastPrimitivesCoderV2_v2_69_9",

                        "ref_Coder_IterableCoder_11",

                    ],

                },



Both its key and value are pickled coded: 



"ref_Coder_DeterministicFastPrimitivesCoderV2_v2_69_9": Coder {

                    spec: Some(

                        FunctionSpec {

                            urn: "beam:coder:pickled_python:v1",



"ref_Coder_IterableCoder_11": Coder {

                    spec: Some(

                        FunctionSpec {

                            urn: "beam:coder:iterable:v1",

                            payload: [],

                        },

                    ),

                    component_coder_ids: [

                        "ref_Coder_FastPrimitivesCoder_6",

                    ],

                },



"ref_Coder_FastPrimitivesCoder_6": Coder {

                    spec: Some(

                        FunctionSpec {

                            urn: "beam:coder:pickled_python:v1",



Whereas java pipeline proto would have beam primitive coder as its components, 
making it straightforward to deserialize and group, how are portable runners 
expected to handle raw pickle kv bytes



Thanks,

Ganesh.  

Reply via email to