You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PyFlink创建bytes类型DataStream报错问题咨询

问题:PyFlink创建bytes类型DataStream报错

我尝试在PyFlink中创建bytes类型的DataStream,执行以下代码时出现错误:

#!/bin/python
# -*- coding: utf-8 -*-

from pyflink.common import Types
from pyflink.datastream import StreamExecutionEnvironment

# create Stream environment for execution (DAG)
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)

message = 'Python is fun'
# convert string to bytes
byte_message = bytes(message, 'utf-8')

ds = env.from_collection(collection=[byte_message], type_info=Types.PRIMITIVE_ARRAY(Types.BYTE())).print()

env.execute()

错误信息

ds = env.from_collection([byte_message], type_info=Types.PRIMITIVE_ARRAY(Types.BYTE())).print()
  File "/home/user/Flink/Python_Projects/pyflink20/lib/python3.10/site-packages/pyflink/datastream/stream_execution_environment.py", line 1009, in from_collection
    return self._from_collection(collection, type_info)
  File "/home/user/Flink/Python_Projects/pyflink20/lib/python3.10/site-packages/pyflink/datastream/stream_execution_environment.py", line 1026, in _from_collection
    j_objs = gateway.jvm.PythonBridgeUtils.readPythonObjects(temp_file.name)
  File "/home/user/Flink/Python_Projects/pyflink20/lib/python3.10/site-packages/py4j/java_gateway.py", line 1322, in __call__
    return_value = get_return_value(
  File "/home/user/Flink/Python_Projects/pyflink20/lib/python3.10/site-packages/pyflink/util/exceptions.py", line 146, in deco
    return f(*a, **kw)
  File "/home/user/Flink/Python_Projects/pyflink20/lib/python3.10/site-packages/py4j/protocol.py", line 326, in get_return_value
    raise Py4JJavaError(
py4j.protocol.Py4JJavaError: An error occurred while calling z:org.apache.flink.api.common.python.PythonBridgeUtils.readPythonObjects.
: java.lang.ClassCastException: class [B cannot be cast to class [Ljava.lang.Object; ([B and [Ljava.lang.Object; are in module java.base of loader 'bootstrap')
at org.apache.flink.api.common.python.PythonBridgeUtils.getObjectArrayFromUnpickledData(PythonBridgeUtils.java:83)
    at org.apache.flink.api.common.python.PythonBridgeUtils.lambda$readPythonObjects$0(PythonBridgeUtils.java:125)
    at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:197)
    at java.base/java.util.LinkedList$LLSpliterator.forEachRemaining(LinkedList.java:1242)
    at java.base/java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:509)
    at java.base/java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:499)
    at java.base/java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:921)
    at java.base/java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234)
    at java.base/java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:682)
    at org.apache.flink.api.common.python.PythonBridgeUtils.readPythonObjects(PythonBridgeUtils.java:131)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:569)
    at org.apache.flink.api.python.shaded.py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
    at org.apache.flink.api.python.shaded.py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374)
    at org.apache.flink.api.python.shaded.py4j.Gateway.invoke(Gateway.java:282)
    at org.apache.flink.api.python.shaded.py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
    at org.apache.flink.api.python.shaded.py4j.commands.CallCommand.execute(CallCommand.java:79)
    at org.apache.flink.api.python.shaded.py4j.GatewayConnection.run(GatewayConnection.java:238)
    at java.base/java.lang.Thread.run(Thread.java:840)

官方文档类型映射说明

Flink 1.20版本文档中明确了以下类型映射关系:

PyFlink Array TypePython TypeJava Type
Types.PRIMITIVE_ARRAY(Types.BYTE())bytesbyte[]

使用环境:Python 3.10.13、Java 17

问题原因与解决方案

这是PyFlink在from_collection方法中处理PRIMITIVE_ARRAY类型时的兼容性问题,底层PythonBridgeUtils会错误地尝试将byte[](原始字节数组)强制转换为Object[],从而触发类型转换异常。

解决方法是使用PyFlink提供的Types.BYTE_ARRAY()类型替代Types.PRIMITIVE_ARRAY(Types.BYTE()),两者的类型映射效果一致,但能避免该转换错误。修改后的代码如下:

#!/bin/python
# -*- coding: utf-8 -*-

from pyflink.common import Types
from pyflink.datastream import StreamExecutionEnvironment

env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)

message = 'Python is fun'
byte_message = bytes(message, 'utf-8')

# 使用Types.BYTE_ARRAY()替代Types.PRIMITIVE_ARRAY(Types.BYTE())
ds = env.from_collection(collection=[byte_message], type_info=Types.BYTE_ARRAY()).print()

env.execute()

内容的提问来源于stack exchange,提问作者IP89

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 07:09:57