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

Kafka集成PySpark报错:UnsatisfiedLinkError问题求助

问题:PySpark读取Kafka流时抛出UnsatisfiedLinkError错误

错误信息

ERROR MicroBatchExecution: Query [id = d3a1ed30-d223-4da4-9052-189b103afca8, runId = 70bfaa84-15c9-4c8b-9058-0f9a04ee4dd0] terminated with error
java.lang.UnsatisfiedLinkError: 'boolean org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(java.lang.String, int)'

生产者代码

import tweepy
import json
import time
from kafka import KafkaProducer
import twitterauth as auth
import utils

producer = KafkaProducer(bootstrap_servers=["localhost:9092"], value_serializer=utils.json_serializer)

class twitterStream(tweepy.StreamingClient):
    def on_connect(self):
        print("Twitter Client Connected")

    def on_tweet(self, raw_data):
        if raw_data.referenced_tweets == None:
            producer.send(topic="registered_user", value=raw_data.text) 
            print("Producer Running")          
            

    def on_error(self):
        self.disconnect()

    def adding_rules(self, keywords):
        for terms in keywords:
            self.add_rules(tweepy.StreamRule(terms))

if __name__ == "__main__":
    stream = twitterStream(bearer_token=auth.bearer_token)
    stream_terms = ['bitcoin','luna','etherum']
    stream.adding_rules(stream_terms)
    stream.filter(tweet_fields=['referenced_tweets'])

PySpark代码

from pyspark.sql import SparkSession
from pyspark import SparkContext
from pyspark.streaming import StreamingContext

import findspark

import json
import os

os.environ['PYSPARK_SUBMIT_ARGS'] = f'--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 pyspark-shell'

if __name__ == "__main__":

    findspark.init()

    sc = SparkSession.builder.master("local[*]") \
                    .appName('SparkByExamples.com') \
                    .getOrCreate()

    df = sc \
        .readStream \
        .format("kafka") \
        .option("kafka.bootstrap.servers", "localhost:9092") \
        .option("startingoffsets","latest") \
        .option("subscribe", "registered_user") \
        .load()
    
    query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
    
    query = df.writeStream.format("console").start()
    import time
    time.sleep(10) # sleep 10 seconds
    query.stop()

环境版本

  • Hadoop版本:3.3.0
  • Spark和PySpark版本:3.3.0
  • Scala版本:2.12.15

补充说明

执行df.printSchema()无报错,报错完整堆栈与Schema输出见附图。


解决方案

该错误是Windows环境下Hadoop原生IO库缺失或权限限制导致的,可尝试以下三种解决方法:

  1. 配置Hadoop winutils工具

    • 下载与Hadoop 3.3.0版本匹配的winutils工具包
    • 设置系统环境变量HADOOP_HOME指向winutils所在目录
    • 将HADOOP_HOME/bin路径添加到系统PATH环境变量中
    • 确认winutils.exe文件存在于HADOOP_HOME/bin目录下
  2. 禁用Hadoop原生IO库
    在构建SparkSession时添加配置,强制使用非原生IO实现:

    sc = SparkSession.builder.master("local[*]") \
                    .appName('SparkByExamples.com') \
                    .config("spark.hadoop.io.native.lib.available", "false") \
                    .getOrCreate()
    
  3. 以管理员身份运行脚本
    Windows系统权限限制可能导致无法访问必要目录,右键点击Python脚本或终端,选择"以管理员身份运行"后再执行代码


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 14:55:20