【发布时间】:2016-01-07 14:17:24
【问题描述】:
我正在尝试将水槽与火花流应用程序集成。我正在运行 spark Scala 示例 FlumePollingEventCount 从水槽中提取事件。我在单机上运行 spark 作业。
我有以下配置。
Avro 源 -> 内存通道 -> Spark SInk
a1.sources = r1
a1.sinks = k1
a1.channels = c1
a1.sources.r1.type = avro
a1.sources.r1.bind = 192.168.1.36
a1.sources.r1.port = 41414
a1.sinks.k1.type = org.apache.spark.streaming.flume.sink.SparkSink
a1.sinks.k1.hostname = 192.168.1.36
a1.sinks.k1.port = 41415
a1.channels.c1.type = memory
a1.channels.c1.capacity = 1000
a1.channels.c1.transactionCapacity = 100
a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1
我正在使用 avro 客户端在 41414 发送事件,但 spark 流无法接收任何事件。
启动 spark 示例时出现以下错误
WARN FlumeBatchFetcher:由于 Flume 代理上的错误,未收到来自 Flume 代理的事件:事务打开时调用 begin()!
在水槽控制台我得到以下异常; 2016-01-07 19:56:51,344(Spark Sink 处理器线程 - 10)[WARN - org.apache.spark.streaming.flume.sink.Logging$class.logWarning(Logging.scala:59)] Spark 无法成功处理事件。事务正在回滚。 2016-01-07 19:56:51,344(新 I/O 工作者 #5)[WARN - org.apache.spark.streaming.flume.sink.Logging$class.logWarning(Logging.scala:59)] 收到错误批处理 - 没有从频道收到任何事件! 2016-01-07 19:56:51,353(新 I/O 工作者 #5)[WARN - org.apache.spark.streaming.flume.sink.Logging$class.logWarning(Logging.scala:59)] 收到错误批处理 - 没有从频道收到任何事件! 2016-01-07 19:56:51,355(Spark Sink 处理器线程 - 9)[WARN - org.apache.spark.streaming.flume.sink.Logging$class.logWarning(Logging.scala:80)] 处理事务时出错. java.lang.IllegalStateException:事务打开时调用 begin()! 在 com.google.common.base.Preconditions.checkState(Preconditions.java:145) 在 org.apache.flume.channel.BasicTransactionSemantics.begin(BasicTransactionSemantics.java:131) 在 org.apache.spark.streaming.flume.sink.TransactionProcessor$$anonfun$populateEvents$1.apply(TransactionProcessor.scala:114) 在 org.apache.spark.streaming.flume.sink.TransactionProcessor$$anonfun$populateEvents$1.apply(TransactionProcessor.scala:113) 在 scala.Option.foreach(Option.scala:236) 在 org.apache.spark.streaming.flume.sink.TransactionProcessor.populateEvents(TransactionProcessor.scala:113) 在 org.apache.spark.streaming.flume.sink.TransactionProcessor.call(TransactionProcessor.scala:243) 在 org.apache.spark.streaming.flume.sink.TransactionProcessor.call(TransactionProcessor.scala:43)
谁能给我一个线索?
【问题讨论】:
标签: apache-spark spark-streaming flume-ng