【问题标题】:Unable to pull events in spark streaming application from flume无法从水槽中提取火花流应用程序中的事件
【发布时间】: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


    【解决方案1】:

    我在使用spark-flume-approach2 时遇到了同样的问题,但在水槽类路径中包含了不同版本的spark-streaming-flume_${spark.scala.version}。如果您包含上述链接中指定的确切版本,则不应再次看到此错误。

    【讨论】:

      【解决方案2】:

      我遇到了同样的问题,这确实是由不同的jar版本引起的。 将 scala-library-2.11.7.jar 替换为 scala-library-2.11.8.jar 后,该问题已得到解决。但是初始消息'begin() 在事务打开时调用!'应该更有意义。 感谢 Sandeep 和 user2710368

      【讨论】:

      • 您能分享一下您在哪里替换了 2.11.7 jar - 在 Flume 中还是在您的 Spark 应用程序中?
      【解决方案3】:

      就我而言,这与 SandeepKumar 非常相似, 在我的 lib 目录中,我有 2 个版本的 scala-Library, 删除旧的,保留所需的解决了我的案子,这花了我 5 个小时。

      【讨论】:

      • 嗨 @user2710368 - Flume 默认在 /lib 中附带 scala-library-2.10.5.jar。而 Spark 需要 scala-library-2.11.8.jar。如果我从水槽的 lib 目录中删除旧的 2.10.5 jar,它会引发异常。你是如何让它与新的 jar 一起工作的?
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-08-09
      • 2015-08-23
      • 2017-04-16
      • 2016-09-29
      • 2023-03-18
      • 1970-01-01
      相关资源
      最近更新 更多