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库缺失或权限限制导致的,可尝试以下三种解决方法:
配置Hadoop winutils工具
- 下载与Hadoop 3.3.0版本匹配的winutils工具包
- 设置系统环境变量
HADOOP_HOME指向winutils所在目录 - 将
HADOOP_HOME/bin路径添加到系统PATH环境变量中 - 确认winutils.exe文件存在于
HADOOP_HOME/bin目录下
禁用Hadoop原生IO库
在构建SparkSession时添加配置,强制使用非原生IO实现:sc = SparkSession.builder.master("local[*]") \ .appName('SparkByExamples.com') \ .config("spark.hadoop.io.native.lib.available", "false") \ .getOrCreate()以管理员身份运行脚本
Windows系统权限限制可能导致无法访问必要目录,右键点击Python脚本或终端,选择"以管理员身份运行"后再执行代码
内容的提问来源于stack exchange,提问作者Bernard Joshua
相关产品推荐
相关产品推荐

