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

使用Java的nats-spark-connector连接NATS JetStream报UnsatisfiedLinkError求助

问题:Spark连接NATS JetStream时触发Hadoop NativeIO UnsatisfiedLinkError错误

场景与代码

使用nats-spark-connector的负载均衡版本编写Java Spark代码,连接NATS JetStream消费消息,核心代码如下:

private static void sparkNatsTester() {
    SparkSession spark = SparkSession.builder()
            .appName("spark-with-nats")
            .master("local")
//          .config("spark.logConf", "false")
              .config("spark.jars",
              "libs/nats-spark-connector-balanced_2.12-1.1.4.jar,"+"libs/jnats-2.17.1.jar"
              )
//            .config("spark.executor.instances", "2")
//            .config("spark.cores.max", "4")
//            .config("spark.executor.memory", "2g")
              .getOrCreate();
    System.out.println("sparkSession : "+ spark);
    Dataset<Row> df = spark.readStream()
            .format("nats")
            .option("nats.host", "localhost")
            .option("nats.port", 4222)
            .option("nats.stream.name", "my_stream")
            .option("nats.stream.subjects", "my_sub")
            // wait 90 seconds for an ack before resending a message
            .option("nats.msg.ack.wait.secs", 1)
            //.option("nats.num.listeners", 2)
            // Each listener will fetch 10 messages at a time
           // .option("nats.msg.fetch.batch.size", 10)
            .load();
    System.out.println("Successfully read nats stream !");
    
    StreamingQuery query;
    try {
        query = df.writeStream()
                  .outputMode("append")
                  .format("console")
                  .start();
        query.awaitTermination();
    } catch (Exception e) {
        e.printStackTrace();
    } 
} 

运行现象与异常

程序能正常打印SparkSession对象和Successfully read nats stream !,随后输出:

Successfully read nats stream !
Status change nats: connection opened
Status change nats: connection closed

随即抛出异常,根因为:

Caused by: java.lang.UnsatisfiedLinkError: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z

完整异常信息:

Exception in thread "stream execution thread for [id = 3ac2d1ac-4876-4c2a-a501-9f94e7e11300, runId = f72897c4-180d-4272-abe2-df9f3838e54b]" org.apache.spark.sql.streaming.StreamingQueryException: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z
=== Streaming Query ===
Identifier: [id = 3ac2d1ac-4876-4c2a-a501-9f94e7e11300, runId = f72897c4-180d-4272-abe2-df9f3838e54b]
Current Committed Offsets: {}
Current Available Offsets: {}

Current State: ACTIVE
Thread State: RUNNABLE

Logical Plan:
WriteToMicroBatchDataSource org.apache.spark.sql.execution.streaming.ConsoleTable$@16cfe41b, 3ac2d1ac-4876-4c2a-a501-9f94e7e11300, Append
+- StreamingExecutionRelation natsconnector.spark.NatsStreamingSource@4b0f9a63, [subject#3, dateTime#4, content#5]

    at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:332)
    at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.$anonfun$run$1(StreamExecution.scala:211)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
    at org.apache.spark.JobArtifactSet$.withActiveJobArtifactState(JobArtifactSet.scala:94)
    at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:211)
Caused by: java.lang.UnsatisfiedLinkError: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z
    at org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Native Method)
    at org.apache.hadoop.io.nativeio.NativeIO$Windows.access(NativeIO.java:793)
    at org.apache.hadoop.fs.FileUtil.canRead(FileUtil.java:1249)
    at org.apache.hadoop.fs.FileUtil.list(FileUtil.java:1454)
    at org.apache.hadoop.fs.RawLocalFileSystem.listStatus(RawLocalFileSystem.java:601)
    at org.apache.hadoop.fs.DelegateToFileSystem.listStatus(DelegateToFileSystem.java:177)
    at org.apache.hadoop.fs.ChecksumFs.listStatus(ChecksumFs.java:548)
    at org.apache.hadoop.fs.FileContext$Util$1.next(FileContext.java:1915)
    at org.apache.hadoop.fs.FileContext$Util$1.next(FileContext.java:1911)
    at org.apache.hadoop.fs.FSLinkResolver.resolve(FSLinkResolver.java:90)
    at org.apache.hadoop.fs.FileContext$Util.listStatus(FileContext.java:1917)
    at org.apache.hadoop.fs.FileContext$Util.listStatus(FileContext.java:1876)
    at org.apache.hadoop.fs.FileContext$Util.listStatus(FileContext.java:1835)
    at org.apache.spark.sql.execution.streaming.AbstractFileContextBasedCheckpointFileManager.list(CheckpointFileManager.scala:315)
    at org.apache.spark.sql.execution.streaming.HDFSMetadataLog.listBatches(HDFSMetadataLog.scala:327)
    at org.apache.spark.sql.execution.streaming.HDFSMetadataLog.getLatest(HDFSMetadataLog.scala:265)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:253)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:427)
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:425)
    at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:67)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$1(MicroBatchExecution.scala:249)
    at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:67)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:239)
    at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:311)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775)
    at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:289)
    ... 4 more

已配置信息

Maven依赖片段

<properties>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <java.version>1.8</java.version>
    <scala.version>2.12</scala.version>
    <spark.version>3.5.0</spark.version>
</properties>

<dependencies>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-core_${scala.version}</artifactId>
        <version>${spark.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-sql_${scala.version}</artifactId>
        <version>${spark.version}</version>
    </dependency>
</dependencies>

系统配置

  • 代码中提前设置:System.setProperty("hadoop.home.dir", "C:\Program Files\Hadoop\winutils-master\hadoop-3.3.1\");
  • 已设置HADOOP_HOME环境变量,并将其bin目录添加至系统Path
  • Maven中Hadoop相关jar版本为3.3.4,曾尝试匹配winutils版本但未解决

修复方案

  • 严格匹配winutils与Hadoop依赖版本:当前Hadoop jar是3.3.4,但winutils用的是3.3.1,版本不兼容导致本地库调用失败。需下载对应3.3.4版本的winutils,替换hadoop.home.dir指向的路径,确保bin目录下包含hadoop.dll、winutils.exe等文件。
  • 禁用Hadoop本地库:在SparkSession配置中添加config("spark.hadoop.io.native.lib.available", "false"),强制Hadoop使用纯Java实现,绕开本地库依赖,这是最快捷的临时解决方案。
  • 检查系统架构一致性:确保winutils版本与JDK架构一致(均为64位或32位),若JDK是64位,必须使用64位的winutils包。
  • 验证文件完整性与权限:检查winutils的bin目录下文件是否完整,且未被杀毒软件隔离;确保当前用户对该目录有读写权限。
  • 重启运行环境:修改hadoop.home.dir或环境变量后,必须重启IDE或终端,确保配置生效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 12:25:55