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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 13:23:30