InfluxDB v3的 Java客户端在Flink的 RichSinkFunction中失败,错误信息为“NameResolver的地址类型 'unix' 不被传输层支持。”
我在一个Apache Flink作业中使用InfluxDB v3 Java客户端。该客户端在独立的Java应用程序中工作正常,但在Flink的 RichSinkFunction中初始化时失败。
Environment
- Apache Flink 1.18.1
- Java 17
- InfluxDB v3 Java client (内部使用gRPC / Arrow Flight)
- Red Hat虚拟机
Minimal example
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) {
// no-op
}
}
用法如下:
stream.addSink(new InfluxSink());
Problem
在 open() 中初始化客户端失败,抛出以下异常:
java.lang.IllegalArgumentException: Address types of NameResolver 'unix' not supported by transport
这在实际写入之前就发生了。
Full Stack Trace
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) [?:?]
What I already verified
- 相同的InfluxDB客户端配置在同一台机器上的独立
main()应用程序中可以正常工作 - 只有在InfluxDB客户端在Flink内部创建时才会出现错误
- 不带scheme的情况下使用
host:8181会导致不同的错误:
java.lang.IllegalArgumentException: java.net.URISyntaxException: Expected scheme-specific part at index 9: grpc+tcp:
- 使用
https://localhost:8181解决了URI解析问题,但随后触发了上文提到的NameResolver 'unix' 错误 - 并行度为
1
Question
是什么原因导致gRPC / Arrow Flight在 Flink内部解析为unix传输?应该如何解决?
Additional Information
io.grpc:grpc 依赖于InfluxDB Java Client
\- 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
解决方案
Solved it — 问题其实是由多种独立问题共同作用导致的,并非Flink/InfluxDB不兼容。
Root Causes
- gRPC / Arrow Flight ServiceLoader在混淆的Fat-JAR中被破坏
InfluxDB v3 Java客户端内部使用gRPC。
我打包的混淆JAR未能把
META-INF/services合并在一起,导致:
Address types of NameResolver 'unix' for 'localhost:8181' not supported by transport
Fix: 将 ServicesResourceTransformer 添加到 maven-shade-plugin。
2. 签名依赖元数据破坏了混淆的JAR
混淆后,签名文件导致:
Invalid signature file digest for Manifest main attributes
Fix: 排除:
META-INF/*.SFMETA-INF/*.DSAMETA-INF/*.RSA- Flink/Kryo + Java 17模块访问问题 除非以以下参数启动,否则Flink序列化会失败:
--add-opens=java.base/java.util=ALL-UNNAMED
4. 本地Docker InfluxDB数据卷权限问题
本地Docker挂载的InfluxDB数据目录写权限错误,导致WAL写入失败:
Operation not permitted (os error 1)
这与Flink/客户端无关,仅影响本地Docker设置。
Result
在修复上述所有问题后,InfluxDB v3 Java客户端在Flink RichSinkFunction 中工作正常,无论是在本地还是目标VM上。
因此,最初的gRPC错误其实是由混淆JAR的打包引起的,而非Flink本身的问题。