【发布时间】:2017-06-07 17:57:01
【问题描述】:
我有一个读取流来使用来自 Kafka 主题的数据,并且基于每个传入消息中的属性值,我必须将数据写入 S3 中的 2 个不同位置中的任何一个(如果 value1 写入 location1,否则位置2)。
在下面的高层次上是我这样做的,
Dataset<Row> kafkaStreamSet = sparkSession
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", kafkaBootstrap)
.option("subscribe", kafkaTopic)
.option("startingOffsets", "latest")
.option("failOnDataLoss", false)
.option("maxOffsetsPerTrigger", offsetsPerTrigger)
.load();
//raw message to ClickStream
Dataset<ClickStream> ds1 = kafkaStreamSet.mapPartitions(processClickStreamMessages, Encoders.bean(ClickStream.class));
ClickStream.java 中有 2 个子对象,一次只会填充其中一个,具体取决于消息属性值是 value1 还是 value2,
1) BookingRequest.java 如果值为 1,
2) PropertyPageView.java if value2 ,
然后我将其从点击流中分离出来,以写入 S3 中的 2 个差异位置,
//fetch BookingRequests in the ClickStream
Dataset<BookingRequest> ds2 = ds1.map(filterBookingRequests,Encoders.bean(BookingRequest.class));
//fetch PropertyPageViews in the ClickStream
Dataset<PropertyPageView> ds3 = ds1.map(filterPropertyPageViews,Encoders.bean(PropertyPageView.class));
最后 ds2 和 ds3 被写入 2 个不同的位置,
StreamingQuery bookingRequestsParquetStreamWriter = ds2.writeStream().outputMode("append")
.format("parquet")
.trigger(ProcessingTime.create(bookingRequestProcessingTime, TimeUnit.MILLISECONDS))
.option("checkpointLocation", "s3://" + s3Bucket+ "/checkpoint/bookingRequests")
.partitionBy("eventDate")
.start("s3://" + s3Bucket+ "/" + bookingRequestPath);
StreamingQuery PageViewsParquetStreamWriter = ds3.writeStream().outputMode("append")
.format("parquet")
.trigger(ProcessingTime.create(pageViewProcessingTime, TimeUnit.MILLISECONDS))
.option("checkpointLocation", "s3://" + s3Bucket+ "/checkpoint/PageViews")
.partitionBy("eventDate")
.start("s3://" + s3Bucket+ "/" + pageViewPath);
bookingRequestsParquetStreamWriter.awaitTermination();
PageViewsParquetStreamWriter.awaitTermination();
它似乎工作正常,并且在部署应用程序时,我看到数据写入了不同的路径。但是,每当作业在失败或手动停止和启动时重新启动时,它都会失败并出现以下异常(其中 userSessionEventJoin.global 是我的主题名称),
原因:org.apache.spark.sql.streaming.StreamingQueryException:预期例如{"topicA":{"0":23,"1":-1},"topicB":{"0":-2}},得到 {"userSessionEventJoin.global":{"92":154362528," 101 org.apache.spark.sql.kafka010.JsonUtils$.partitionOffsets(JsonUtils.scala:74) org.apache.spark.sql.kafka010.KafkaSourceOffset$.apply(KafkaSourceOffset.scala:59)
如果我删除了所有的检查点信息,那么它会再次开始并在给定的 2 个位置开始新的检查点,但这意味着我必须再次从最新的偏移量开始处理并丢失所有以前的偏移量。
spark版本是2.1,这个topic有100+个partition。
我只用一个写入流(一个检查点位置)进行了测试,重启时会发生同样的异常。
请提出任何解决方案,谢谢。
【问题讨论】:
标签: apache-spark spark-streaming apache-spark-2.0