PyFlink中出现gRPC连接取消错误:Multiplexer Hanging Up
问题描述
在Mac M2环境下,使用Flink 1.17 + Python 3.9实现官方文档中的Python数据流示例时,触发如下错误:
Traceback (most recent call last): File "/Library/Developer/CommandLineTools/Library/Frameworks/Python3.framework/Versions/3.9/lib/python3.9/threading.py", line 973, in _bootstrap_inner self.run() File "/Library/Developer/CommandLineTools/Library/Frameworks/Python3.framework/Versions/3.9/lib/python3.9/threading.py", line 910, in run self._target(*self._args, **self._kwargs) File "/Users/shashikanth/Downloads/flink-1.17.2/pyflink_env/lib/python3.9/site-packages/apache_beam/runners/worker/data_plane.py", line 669, in <lambda> target=lambda: self._read_inputs(elements_iterator), File "/Users/shashikanth/Downloads/flink-1.17.2/pyflink_env/lib/python3.9/site-packages/apache_beam/runners/worker/data_plane.py", line 652, in _read_inputs for elements in elements_iterator: File "/Users/shashikanth/Downloads/flink-1.17.2/pyflink_env/lib/python3.9/site-packages/grpc/_channel.py", line 542, in __next__ return self._next() File "/Users/shashikanth/Downloads/flink-1.17.2/pyflink_env/lib/python3.9/site-packages/grpc/_channel.py", line 968, in _next raise self grpc._channel._MultiThreadedRendezvous: <_MultiThreadedRendezvous of RPC that terminated with: status = StatusCode.CANCELLED details = "Multiplexer hanging up" debug_error_string = "UNKNOWN:Error received from peer ipv6:%5B::1%5D:50616 {created_time:"2024-03-20T10:56:21.548564+05:30", grpc_status:1, grpc_message:"Multiplexer hanging up"}"
原因分析
- 这是PyFlink底层依赖的Apache Beam与gRPC交互时的常见问题,核心是多路复用连接异常中断导致RPC调用被取消。
- Mac M2的ARM架构与Flink 1.17绑定的旧版本gRPC/Apache Beam存在兼容性适配问题,部分底层逻辑未针对ARM架构优化。
- 本地IPv6配置冲突:错误信息中的
ipv6:%5B::1%5D表明gRPC尝试通过IPv6建立连接时出现异常,直接导致连接终止。
解决建议
- 强制gRPC使用IPv4:运行Flink任务前设置环境变量,规避IPv6连接异常:
export GRPC_PYTHON_FORCE_IPV4=1 export GRPC_DNS_RESOLVER=native - 升级依赖版本:在PyFlink的虚拟环境中更新gRPC和Apache Beam,适配ARM架构:
pip install --upgrade grpcio apache-beam - 调整Flink配置:在
flink-conf.yaml中添加以下参数,优化Python Worker的连接性能:python.flink.grpc.server.num-threads: 4 python.flink.execution-mode: cluster - 验证基础环境:先运行极简的WordCount示例,排除业务代码问题,确认环境可用性:
from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import StreamTableEnvironment env = StreamExecutionEnvironment.get_execution_environment() t_env = StreamTableEnvironment.create(env) # 创建数据源表 t_env.execute_sql(""" CREATE TABLE SourceTable ( word STRING ) WITH ( 'connector' = 'datagen', 'rows-per-second' = '10' ) """) # 创建输出表 t_env.execute_sql(""" CREATE TABLE SinkTable ( word STRING, cnt BIGINT ) WITH ( 'connector' = 'print' ) """) # 执行统计逻辑 t_env.execute_sql(""" INSERT INTO SinkTable SELECT word, COUNT(*) FROM SourceTable GROUP BY word """).wait()
内容的提问来源于stack exchange,提问作者Koder
相关产品推荐
相关产品推荐

