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 Type | Python Type | Java Type |
|---|---|---|
| Types.PRIMITIVE_ARRAY(Types.BYTE()) | bytes | byte[] |
使用环境: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
相关产品推荐
相关产品推荐

