【问题标题】:Why different count results on consecutive reads?为什么连续读取的计数结果不同?
【发布时间】:2017-09-28 23:34:36
【问题描述】:

我有以下代码正在将表读取到 apache spark DataFrame:

val df = spark.read.format("jdbc").option("url", "jdbc:postgresql:host/database").option("dbtable", "tablename").option("user", "username").option("password", "password").load()

当我第一次调用 df.count() 时,我得到的数字比下次调用相同的 count 方法时要小。

为什么会这样?

当我第一次读取该表时,Spark 不会在我的 Spark 集群上的 DataFrame 中加载我的表的快照吗?

我在 postgres 上的表格不断被喂食,我的数据框似乎反映了这种行为。

在调用read 方法时,我应该如何设法仅将我的表的静态快照加载到 spark 的 DataFrame 中?

【问题讨论】:

  • 你试过在 postgresql 端检查表行数吗?
  • @RamGhadiyaram 是的,这与我从 spark 端的 count 方法得到的结果完全相同。所以我的 spark DataFrame 确实反映了我在 postgres 上的表中的这种动态行为
  • 我认为...可能是 spark 只读取 postgre sql 的已提交数据,这就是为什么您的 df.count 不时变化的原因。
  • @RamGhadiyaram,是的,这听起来像是在这里发生的事情。我认为 Spark 读取所有内容一次,然后对 DataFrame 的后续调用会考虑加载到 Spark Cluster 上的我的表的静态副本。但似乎并非如此。
  • "在调用 read 方法时,我应该如何设法仅将我的表的静态快照加载到 spark 的 DataFrame 中?" - 你不能在postgre表级别引入一个快照id,snapshot_datetime,然后根据这些参数从spark sql中读取它吗?

标签: postgresql apache-spark jdbc apache-spark-sql


【解决方案1】:

除非Datasetcached 使用可靠存储(标准Spark cache 只会给你弱保证)数据库可能会被多次访问,每次都显示数据库的当前状态。自从

postgres 上的表格不断被喂食

看到不同的计数是一种预期的行为。

此外,如果 JDBC 源用于分布式模式(带有分区列或predicates),那么每个执行器线程将使用自己的事务。因此Dataset 的状态可能不完全一致。

我应该如何设法只加载静态快照

不要使用 JDBC。例如,您可以

  • COPY 数据到文件系统并从那里加载。
  • 使用您选择的复制解决方案创建专用于分析的副本,并在分析数据时设置和暂停复制。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-02-11
    • 2012-08-09
    相关资源
    最近更新 更多