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

Flink RichSinkFunction中InfluxDB v3 Java客户端初始化遇unix解析错误

Flink中InfluxDB v3 Java客户端初始化失败问题

环境

  • Apache Flink 1.18.1
  • Java 17
  • InfluxDB v3 Java客户端(内部基于gRPC/Arrow Flight)
  • RedHat虚拟机

最小复现示例

import com.influxdb.v3.client.InfluxDBClient;
import com.influxdb.v3.client.config.ClientConfig;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;

public class InfluxSink extends RichSinkFunction<String> {

    private transient InfluxDBClient client;

    @Override
    public void open(Configuration parameters) {
        ClientConfig config = new ClientConfig.Builder()
                .host("https://host:8181")
                .token("TOKEN".toCharArray())
                .database("DB")
                .build();

        this.client = InfluxDBClient.getInstance(config);
    }

    @Override
    public void invoke(String value, Context context) {
        // 空实现
    }
}

使用方式:

stream.addSink(new InfluxSink());

问题描述

客户端在open()方法初始化阶段抛出如下异常,该错误发生在任何实际写入操作之前:

java.lang.IllegalArgumentException: Address types of NameResolver 'unix' not supported by transport

完整堆栈跟踪

java.lang.IllegalArgumentException: Address types of NameResolver 'unix' for 'myhost.com:8181' not supported by transport
        at io.grpc.internal.ManagedChannelImplBuilder.getNameResolverProvider(ManagedChannelImplBuilder.java:871) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at io.grpc.internal.ManagedChannelImplBuilder.build(ManagedChannelImplBuilder.java:721) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at io.grpc.ForwardingChannelBuilder2.build(ForwardingChannelBuilder2.java:278) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at com.influxdb.v3.client.internal.FlightSqlClient.createFlightClient(FlightSqlClient.java:179) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at com.influxdb.v3.client.internal.FlightSqlClient.<init>(FlightSqlClient.java:102) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at com.influxdb.v3.client.internal.FlightSqlClient.<init>(FlightSqlClient.java:82) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at com.influxdb.v3.client.internal.InfluxDBClientImpl.<init>(InfluxDBClientImpl.java:116) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at com.influxdb.v3.client.internal.InfluxDBClientImpl.<init>(InfluxDBClientImpl.java:97) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at com.influxdb.v3.client.InfluxDBClient.getInstance(InfluxDBClient.java:519) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at sink.InfluxSink.open(InfluxSink.java:46) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.api.common.functions.util.FunctionUtils.openFunction(FunctionUtils.java:34) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.open(AbstractUdfStreamOperator.java:101) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.streaming.api.operators.StreamSink.open(StreamSink.java:46) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.initializeStateAndOpenOperators(RegularOperatorChain.java:107) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreGates(StreamTask.java:753) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.call(StreamTaskActionExecutor.java:55) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreInternal(StreamTask.java:728) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.streaming.runtime.tasks.StreamTask.restore(StreamTask.java:693) ~[ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:953) [ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:922) [ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:746) [ApacheFlink-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.runtime.taskmanager.Task.run(Task.java:562) [ApacheFlink-1.0-SNAPSHOT.jar:?]
        at java.lang.Thread.run(Thread.java:840) [?:?]

已验证信息

  • 同一InfluxDB客户端配置在同一机器的独立main()应用中可正常运行
  • 仅在Flink内部创建InfluxDB客户端时触发该错误
  • 使用不带协议的host:8181会引发URI解析错误:
java.lang.IllegalArgumentException: java.net.URISyntaxException: Expected scheme-specific part at index 9: grpc+tcp:
  • 使用https://localhost:8181可修复URI解析问题,但仍会触发上述NameResolver 'unix'错误
  • 作业并行度已设置为1

疑问

为何gRPC/Arrow Flight在Flink环境中会解析为unix传输?该如何修复此问题?

额外依赖信息

InfluxDB Java客户端的io.grpc:grpc依赖树:

\- com.influxdb:influxdb3-java:jar:1.8.0:compile
   \- org.apache.arrow:flight-core:jar:18.3.0:compile
      +- io.grpc:grpc-netty:jar:1.71.0:compile
      |  \- io.grpc:grpc-util:jar:1.71.0:runtime
      +- io.grpc:grpc-core:jar:1.71.0:compile
      |  \- io.grpc:grpc-context:jar:1.71.0:runtime
      +- io.grpc:grpc-protobuf:jar:1.71.0:compile
      |  \- io.grpc:grpc-protobuf-lite:jar:1.71.0:runtime
      +- io.grpc:grpc-stub:jar:1.71.0:compile
      \- io.grpc:grpc-api:jar:1.71.0:compile

问题分析与解决方案

原因

Flink环境的类加载机制与独立Java应用不同,导致gRPC加载了错误的NameResolverProvider。Flink类路径中可能存在支持unix域套接字的resolver实现,gRPC服务发现机制优先选择了该实现,但当前使用的Netty传输不支持unix地址类型,从而抛出异常。

修复方案

方案1:强制指定gRPC使用DNS NameResolver

在Flink作业启动时添加JVM参数,强制gRPC使用DNS解析器:

-Dgrpc.nameResolverFactory=io.grpc.internal.DnsNameResolverProvider

或者在代码中提前设置系统属性(需确保在客户端初始化前执行):

static {
    System.setProperty("grpc.nameResolverFactory", "io.grpc.internal.DnsNameResolverProvider");
}

方案2:自定义gRPC通道配置

修改InfluxSink的open()方法,显式配置gRPC通道使用DNS解析器:

@Override
public void open(Configuration parameters) {
    ClientConfig config = ClientConfig.builder()
            .host("https://host:8181")
            .token("TOKEN".toCharArray())
            .database("DB")
            .build();

    // 自定义FlightClient的gRPC通道构建器
    FlightClient.Builder flightBuilder = FlightClient.builder();
    flightBuilder.grpcChannelBuilder(
            NettyChannelBuilder.forTarget(config.getHost().replaceFirst("https://", ""))
                    .nameResolverFactory(DnsNameResolverProvider.getInstance())
                    .sslContext(GrpcSslContexts.forClient().build())
    );

    this.client = InfluxDBClient.getInstance(config, flightBuilder);
}

注意:需确保引入io.grpc:grpc-netty和io.netty:netty-tcnative-boringssl-static相关依赖,以支持SSL上下文构建。

方案3:调整依赖排除不必要的NameResolver

检查Flink作业的依赖树,排除可能引入unix NameResolver的依赖。例如在pom.xml中排除相关依赖:

<dependency>
    <groupId>com.influxdb</groupId>
    <artifactId>influxdb3-java</artifactId>
    <version>1.8.0</version>
    <exclusions>
        <exclusion>
            <groupId>io.grpc</groupId>
            <artifactId>grpc-unix-domain-socket</artifactId>
        </exclusion>
    </exclusions>
</dependency>

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 17:54:54