【问题标题】:Apache Spark - capturing Kafka data on streaming event to trigger workflowApache Spark - 在流式事件中捕获 Kafka 数据以触发工作流
【发布时间】:2019-03-16 20:02:25
【问题描述】:

简而言之,我是一名尝试使用 Spark 将数据从一个系统移动到另一个系统的开发人员。一个系统中的原始数据经过处理、汇总后形成一个本土分析系统。

我对 Spark 非常陌生 - 我的知识仅限于我在过去一两周内能够挖掘和试验的内容。

我所描绘的是;使用 Spark 监视来自 Kafka 的事件作为触发器。捕获有关消费者事件的实体/数据,并使用它来告诉我分析系统中需要更新的内容。然后,我将对原始 Cassandra 数据运行相关的 Spark 查询,并将结果写入分析端的另一个表中,仪表板指标将其称为数据源。

我有一个简单的 Kafka 结构化流式查询工作。虽然我可以看到消费的对象正在输出到控制台,但当消费者事件发生时,我无法检索 Kafka 记录:

try {
    SparkSession spark = SparkSession
        .builder()
        .master(this.sparkMasterAddress)
        .appName("StreamingTest2")
        .getOrCreate();

    //THIS -> None of these events seem to give me the data consumed?
    //...thinking I'd trigger the Cassandra write from here?
    spark.streams().addListener(new StreamingQueryListener() {
        @Override
        public void onQueryStarted(QueryStartedEvent queryStarted) {
            System.out.println("Query started: " + queryStarted.id());
        }
        @Override
        public void onQueryTerminated(QueryTerminatedEvent queryTerminated) {
            System.out.println("Query terminated: " + queryTerminated.id());
        }
        @Override
        public void onQueryProgress(QueryProgressEvent queryProgress) {
            System.out.println("Query made progress: " + queryProgress.progress());
        }
    });

    Dataset<Row> reader = spark
        .readStream()
        .format("kafka")
        .option("startingOffsets", "latest")
        .option("kafka.bootstrap.servers", "...etc...")
        .option("subscribe", "my_topic")
        .load();

    Dataset<String> lines = reader
        .selectExpr("cast(value as string)")
        .as(Encoders.STRING());

    StreamingQuery query = lines
        .writeStream()
        .format("console")
        .start();
    query.awaitTermination();
} catch (Exception e) {
    e.printStackTrace();
}

我还可以使用 Spark SQL 很好地查询 Cassandra:

try {
    SparkSession spark = SparkSession.builder()
        .appName("SparkSqlCassandraTest")
        .master("local[2]")
        .getOrCreate();

    Dataset<Row> reader = spark
        .read()
        .format("org.apache.spark.sql.cassandra")
        .option("host", this.cassandraAddress)
        .option("port", this.cassandraPort)
        .option("keyspace", "my_keyspace")
        .option("table", "my_table")
        .load();

    reader.printSchema();
    reader.show();

    spark.stop();
} catch (Exception e) {
    e.printStackTrace();
}

我的想法是;用前者触发后者,将这个东西捆绑为 Spark 应用程序/包/任何东西,并将其部署到 spark 中。到那时,我希望它会不断将更新推送到指标表。

这是否是满足我需要的可行、可扩展、合理的解决方案?我在正确的道路上吗?不反对在某种程度上使用更容易或更好的 Scala。

谢谢!

编辑:这是我所反对的图表。

【问题讨论】:

  • FWIW,假设 Spark 实际上没有转换任何数据,那么 Kafka Connect 旨在读取 Kafka 主题并将这些事件写入下游系统,而无需自己编写任何代码(Cassandra 连接器已经存在) .这样,您无需弄清楚如何在本地计算机之外部署和监控长时间运行的 Spark Streaming 作业
  • @cricket_007 我需要执行转换以从左侧数据中处理原始数据,并在右侧系统上更新一组表格,这是显示结果的分析层/metrics 来自这些转换。我只想使用 Kafka 作为原始系统中原始数据发生变化的通知。
  • 知道了...我对 Spark 不太熟悉,不知道你是否真的可以逐行访问低级 Kafka Consumer 事件记录。

标签: java scala apache-spark cassandra apache-kafka


【解决方案1】:

知道了。了解了 ForeachWriter。效果很好:

        StreamingQuery query = lines
            .writeStream()
            .format("foreach")
            .foreach(new ForeachWriter<String>() {
                @Override
                public void process(String value) {
                    System.out.println("process() value = " + value);
                }

                @Override
                public void close(Throwable errorOrNull) {}

                @Override
                public boolean open(long partitionId, long version) {
                    return true;
                }
            })
            .start(); 

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-12-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-05-12
    • 1970-01-01
    相关资源
    最近更新 更多