【问题标题】:Is it good practice to query a database when streaming with Spark [closed]使用 Spark 进行流式传输时查询数据库是否是一种好习惯 [关闭]
【发布时间】:2020-07-04 04:29:51
【问题描述】:

我正在尝试使用 Spark(以及 Flink)制作一个简单的流媒体应用程序。输入是一个 Kafka 主题,其中包含一些产品信息,在我将此产品信息写入其他地方之前,我想获取产品的类别信息并在流式传输期间添加到我的模型中。我打算将我的产品类别映射数据存储在 PostgreSQL 表或 Couchbase 存储桶中。所以我需要为每个流数据查询其中一个。我想学的是:

  • 是否适用于 Spark/Flink?

  • 对于流媒体应用来说这是一个好习惯吗?如果没有,我怎样才能用大数据的方式做到这一点?我应该将我的地图数据存储在其他地方,还是以其他方式将该地图数据与流数据结合起来?

【问题讨论】:

  • 为什么要将映射数据存储在 PostgreSQL 或 Couchbase 中?,您可以将相同的数据推送到 kafka 并编写一个批处理作业以将相同的 kafka 数据同步到数据库??
  • 每天都会有新的条目出现在产品类别映射表中。那么,我是否应该每天将映射表的所有条目推送到 Kafka,并在 Flink 上将这个新的映射流加入产品流?如果我误解了,请纠正我。
  • Flink 我没工作过.. 但是在 spark 或 kafka 流中你可以做同样的事情......一个原因是为每个记录流查询 rdbms 数据库是昂贵的操作.. 你可能会想到替代品..

标签: apache-spark spark-streaming apache-flink flink-streaming


【解决方案1】:

这至少对于 Spark 来说是相当普遍的做法,但它可能更多地取决于数据库选择。就像 Cassandra 在通过完整主键进行查找时非常快(如果您没有大量数据,启用行缓存也可能会有所帮助),我认为 Couchbase 也可能是一个不错的选择(我没有使用它需很长时间)。代码可能看起来像这样(完整代码在 this Zeppelin Notebook 中。此代码还需要 Spark Cassandra 连接器 2.5.0,它支持“直接加入 Cassandra” - 请参阅 blog post on that release):

val streamingInputDF = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "10.101.36.9:9092")
      .option("subscribe", "tweets2")
      .load()

val tweetDF = streamingInputDF.selectExpr("CAST(value AS STRING)")
  .select(from_json($"value", schema).as("tweet"))
  .select($"tweet.payload.created_at".as("created_at").cast(TimestampType),
    $"tweet.payload.lang".as("language"))

val streamingCountsDF =  tweetDF
  .where(col("language").isNotNull)
  .groupBy($"language", window($"created_at", "1 minutes"))
  .count()
  .select($"language", $"window.start".as("ts"), $"count")

import org.apache.spark.sql.cassandra._
val lang_details = spark.read.cassandraFormat("languages", "zep").load()

val joined = streamingCountsDF.join(lang_details, 
       lang_details("id") === streamingCountsDF("language"), "left_outer")
  .select($"language", $"native_name".as("lang_name"), $"ts", $"count")

...

此外,您还需要考虑其他要求 - 您获得更新的频率、这些更新需要以多快的速度传播到流式作业等。例如,对于 Spark,您可能有一个单独的数据框用于存储在数据库中的数据,您可以缓存这些数据以便更快地加入,但每 N 分钟刷新一次数据帧,以从 DB 获取最新更新。 (你可以在 Stream Processing with Apache Spark 书中找到源代码)

【讨论】:

    【解决方案2】:

    这是 Flink 应用程序中的常见模式,有几种方法可以考虑:

    (1) 您可以使用外部数据库进行查找连接。 Table API 内置了对 JDBC 数据库(包括 PostgreSQL)的支持。例如here is an enrichment join of a Kafka stream with a lookup table in MySQL,通过 Hive 目录访问 MySQL 表:

    SELECT
      l_proctime AS `querytime`,
      l_orderkey AS `order`,
      l_linenumber AS `linenumber`,
      l_currency AS `currency`,
      rs_rate AS `cur_rate`, 
      (l_extendedprice * (1 - l_discount) * (1 + l_tax)) / rs_rate AS `open_in_euro`
    FROM prod_lineitem
    JOIN hive.`default`.prod_rates FOR SYSTEM_TIME AS OF l_proctime ON rs_symbol = l_currency
    WHERE
      l_linestatus = 'O';
    

    Documentation.

    (2) 对于其他(非 JDBC)数据源,您可以使用 async function 使用外部服务/数据库实现自己的扩充。

    (3) Flink 1.11 增加了对摄取 Debezium CDC(变更数据捕获)流的支持,从而更容易维护以 Flink 状态物化的外部数据库的同步视图。支持 MySQL、PostgreSQL、Oracle、Microsoft SQL Server 和许多其他数据库。通过将外部数据源镜像到 Flink 中,您将获得更高的吞吐量和更低的延迟,并减少外部数据库的负载。

    【讨论】:

      猜你喜欢
      • 2015-02-18
      • 2013-03-30
      • 1970-01-01
      • 1970-01-01
      • 2016-03-20
      • 1970-01-01
      • 2010-09-14
      • 2020-09-10
      • 2017-04-28
      相关资源
      最近更新 更多