【问题标题】:How to link Spark output to Logstash input如何将 Spark 输出链接到 Logstash 输入
【发布时间】:2016-11-28 13:34:08
【问题描述】:

我有一个 Spark Streaming 作业输出一些当前存储在 HDFS 中的日志,我想用 logstash 处理它们。不幸的是,虽然有一个插件可以在 hdfs 中为 logstash 编写,但实际上用它从 hdfs 中读取是不可能的。

我已经搜索了链接这两个部分的解决方案,但就 python api 的 Spark 流而言,存储内容的唯一方法是将其作为文本文件写入 hdfs,所以我必须从 hdfs 读取! 我无法将它们保存在本地,因为 Spark 在集群上运行,并且我不想从每个节点获取所有数据。

目前我运行一个非常脏的脚本,它每 2 秒将内容复制到本地 hdfs 目录。但是这个解决方案显然不能令人满意。

有人知道可以帮助我将 Spark 的输出发送到 Logstash 的软件吗?

提前致谢!

编辑:我使用 Python 和 Spark 1.6.0

【问题讨论】:

  • 这些是Log4j生成的日志吗?
  • 不,这是 Spark 处理的 apache 日志,它基于机器学习算法为其添加了一些功能。

标签: python apache-spark hdfs logstash spark-streaming


【解决方案1】:

这似乎是使用Kafka 的完美工作。在您的 Spark Streaming 作业中,写入 Kafka,然后使用 Logstash 中的记录。

stream.foreachRDD { rdd =>
  rdd.foreachPartition { partition =>
    val producer = createKafkaProducer()
    partition.foreach { message =>
      val record = ... // convert message to record
      producer.send(record)
    }
    producer.close()
  }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-03-10
    • 1970-01-01
    相关资源
    最近更新 更多