【问题标题】:How to design spark program to process 300 most recent files?如何设计 Spark 程序来处理 300 个最新文件?
【发布时间】:2016-04-12 02:00:27
【问题描述】:

情况

新的小文件会定期出现。我需要对最近的 300 个文件进行计算。所以基本上有一个前进的窗口。窗口大小为 300,需要对窗口进行计算。

但是要知道的非常重要的一点是,这不是火花流计算。因为在火花流中,窗口的单位/范围是时间。这里的单位/范围是文件数。

解决方案1

我会维护一个dict,dict的大小是300。每个新文件进来,我把它变成spark数据框,放入dict。然后我确保 dict 中最旧的文件被弹出 如果 dict 的长度超过 300。 在此之后,我会将 dict 中的所有数据帧合并为一个更大的数据帧并进行计算。

上述过程将循环运行。每次新文件进来时,我们都会循环。

解决方案 1 的伪代码

for file in file_list:
    data_frame = get_data_frame(file)
    my_dict[ timestamp ] = data_frame

    for timestamp in my_dict.keys():
        if timestamp older than 24 hours:
            # not only unpersist, but also delete to make sure the memory is released 
            my_dict[timestamp].unpersist
            del my_dict[ timestamp ] 

    # pop one data frame from the dict
    big_data_frame = my_dict.popitem()

    for timestamp in my_dict.keys():
        df = my_dict.get( timestamp )
        big_data_frame = big_data_frame.unionAll(df)

    # Then we run SQL on the big_data_frame to get report

解决方案 1 的问题

总是达到内存不足gc开销限制

问题

您是否发现解决方案 1 中有任何不当之处?

有没有更好的解决方案?

这是使用 spark 的正确情况吗?

【问题讨论】:

  • 你说你的计算窗口是300对吧?但是在解决方案 1 中,如果您弹出最旧的文件,那么您仍然有 299 个旧文件,对吗?你能澄清我的理解吗?
  • @LokeshKumarP 嗨,我修改了问题。在弹出数据之前,我会检查字典。如果 dict 的总长度没有达到 300,那么我不会弹出任何东西。如果还不清楚,请告诉我
  • 感谢您的澄清,每个文件的大小是多少,而且您提到它不是流式作业,那么您多久轮询一次文件系统?
  • @LokeshKumarP 我添加了一些伪代码。文件大小为 4KB。该集群是火花独立集群。根据频率。一旦循环结束,只要目录中还有文件,我就会执行下一个循环。
  • 当你拥有大数据帧时,它有多少个分区?单个或多个,以及您正在运行什么查询,因为内存使用还取决于您正确使用的 SQL 构造?

标签: apache-spark


【解决方案1】:

一个观察,你可能不想使用popitem,Python 字典的键没有排序,所以你不能保证你弹出最早的项目。相反,我每次都会使用排序的时间戳列表重新创建字典。假设您的文件名只是时间戳:

my_dict = {file:get_dataframe(file) for file in sorted(file_list)[-300:]}

不确定这是否能解决您的问题,您能否将错误的完整堆栈跟踪粘贴到问题中?您的问题可能发生在 Spark 合并/连接中(未包含在您的问题中)。

【讨论】:

  • 感谢 maxymoo。我实际上使用时间戳作为 dict 的键。所以找到最旧的数据框不会有问题。堆栈跟踪太多了。但你是对的,我稍后会发布它。谢谢
  • 不,我的意思是,如果您要使用 pop 删除最旧的项目,则需要使用 OrderedDict - 普通的 dict 并不总是对键进行排序,所以你不能保证pop 会得到最早的日期。
  • 谢谢。我实际上并没有真正使用简单的弹出方法。我使用普通的字典。然后用key判断对应的item是否应该被移除。
【解决方案2】:

我对此的建议是流式传输,但不是关于时间,我的意思是你仍然会设置一些窗口和滑动间隔,但说它是 60 秒。

因此,您每 60 秒就会在“x”分区中获得文件内容的 DStream。这些“x”分区代表您拖放到 HDFS 或文件系统上的文件。 因此,通过这种方式,您可以跟踪已读取的文件/分区数量,如果它们少于 300,则等到它们变为 300。当计数达到 300 后,您就可以开始处理了。

【讨论】:

    【解决方案3】:

    如果可以跟踪最新的文件,或者可以偶尔发现它们,那么我建议做类似的事情

    sc.textFile(','.join(files));
    

    或者如果可以识别特定模式来获取这 300 个文件,那么

    sc.textFile("*pattern*");
    

    甚至可以有逗号分隔的模式,但可能会发生一些匹配多个模式的文件,而不是一次读取。

    【讨论】:

    • 您好,谢谢。是的,我们可以使用模式方式。但是问题就像你说的那样,有些文件可能会被读取多次。这就是为什么我使用 dict 来跟踪那些已经读取的文件。我想避免那些读取成本
    • 然后将它们加入逗号分隔列表并每次阅读。代码就这么简单,所以我不会太担心磁盘读取性能。有 SSD,有磁盘缓冲区。阅读时间与处理时间的百分比是多少。无论如何,如果这是一个问题,我会测量时间来阅读,但很可能最终会得到更简单的代码 =)
    猜你喜欢
    • 2021-11-29
    • 1970-01-01
    • 1970-01-01
    • 2022-07-05
    • 2013-06-04
    • 1970-01-01
    • 1970-01-01
    • 2015-07-29
    • 2015-05-29
    相关资源
    最近更新 更多