【发布时间】:2018-10-06 20:52:31
【问题描述】:
我尝试使用 Spark 从 Kafka 消费,更具体地说是 PySpark 和 Structured Streaming。
import os
import time
import time
from ast import literal_eval
from pyspark.sql.types import *
from pyspark.sql.functions import from_json, col, struct, explode
from pyspark.sql import SparkSession
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.3.0 pyspark-shell'
spark = SparkSession \
.builder \
.appName("Structured Streaming") \
.getOrCreate()
requests = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "ip-ec2:9092") \
.option("subscribe", "ssp.requests") \
.option("startingOffsets", "earliest") \
.load()
requests.printSchema()
# root |-- key: binary (nullable = true) |-- value: binary (nullable =
# true) |-- topic: string (nullable = true) |-- partition: integer
# (nullable = true) |-- offset: long (nullable = true) |-- timestamp:
# timestamp (nullable = true) |-- timestampType: integer (nullable =
# true)
当我运行下一行代码时
rawQuery = requests \
.selectExpr("topic", "CAST(key AS STRING)", "CAST(value AS STRING)") \
.writeStream.trigger(processingTime="5 seconds") \
.format("parquet") \
.option("checkpointLocation", "/home/user/folder/applicationHistory") \
.option("path", "/home/user/folder") \
.start()
rawQuery.awaitTermination()
Py4JJavaError Traceback(最近调用 最后)/opt/conda/lib/python3.6/site-packages/pyspark/sql/utils.py 在 装饰(*a, **kw) 62 尝试: ---> 63 返回 f(*a, **kw) 64 除了 py4j.protocol.Py4JJavaError as e:
/opt/conda/lib/python3.6/site-packages/py4j/protocol.py 在 get_return_value(answer, gateway_client, target_id, name) 319 “调用 {0}{1}{2} 时出错。\n”。 --> 320 格式(target_id, ".", name), value) 321 其他:
Py4JJavaError:调用 o70.awaitTermination 时出错。 : org.apache.spark.sql.streaming.StreamingQueryException:作业中止。 === 流式查询 === 标识符:[id = c2b48840-5ba4-416e-a192-dcae94007856, runId = 4afcca20-00cd-4187-a70b-1b742f1f5c0d] 当前提交的偏移量:{} 当前可用的偏移量:{KafkaSource[Subscribe[ssp.requests]]:
我无法理解这个错误的原因
Py4JJavaError: 调用 o70.awaitTermination 时出错
【问题讨论】:
-
你能发布完整的堆栈跟踪吗(这是 Py4J 异常,所以 JVM 部分很重要)。
标签: apache-spark pyspark apache-kafka spark-structured-streaming