【发布时间】: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