【发布时间】:2017-08-22 18:04:21
【问题描述】:
我想通过 PySpark 使用结构化流式处理运行 Spark 应用程序。
我使用 Spark 2.2 和 Kafka 0.10 版本。
我因以下错误而失败:
java.lang.IncompatibleClassChangeError:实现类
spark-submit 命令使用如下:
/bin/spark-submit \
--packages org.apache.spark:spark-streaming-kafka-0-10_2.11:2.2.0 \
--master local[*] \
/home/umar/structured_streaming.py localhost:2181 fortesting
structured_streaming.py代码:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("StructuredStreaming").config("spark.driver.memory", "2g").config("spark.executor.memory", "2g").getOrCreate()
raw_DF = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "localhost:2181").option("subscribe", "fortesting").load()
values = raw_DF.selectExpr("CAST(value AS STRING)").as[String]
values.writeStream.trigger(ProcessingTime("5 seconds")).outputMode("append").format("console").start().awaitTermination()
【问题讨论】:
标签: apache-spark pyspark spark-structured-streaming