chunk.apply

R 的可扩展数据处理

Simon Urbanek

Member of R-Core, Lead Inventive Scientist, AT&T Labs Research

chunk.apply()

  • 抽象循环过程
  • 支持并行执行
  • iotoolshmr 的基础,可在 Apache Hadoop 基础设施上处理数据
R 的可扩展数据处理

mstrsplit() 将分块读为矩阵

# 使用 chunk.apply 从 foo.csv 获取按行分块
chunk_col_sums <- chunk.apply("foo.csv",

# 处理每个分块的函数 function(chunk) { # 将分块转为矩阵 m <- mstrsplit(chunk, type = "numeric", sep = ",") # 返回列和 colSums(m) }, # 最大分块字节数 CH.MAX.SIZE = 1e5)
# 获取总和 colSums(chunk_col_sums)
R 的可扩展数据处理

dstrsplit() 将分块读为数据框

# 使用 chunk.apply 从 foo.csv 获取按行分块
chunk_col_sums <- chunk.apply("foo.csv",

 # 处理每个分块的函数
 function(chunk) {
   # 将分块转为数据框
   d <- dstrsplit(chunk, col_types = rep("numeric", 3), sep = ",")
   # 返回列和
   colSums(d)
 }, 
 # 最大分块字节数
 CH.MAX.SIZE = 1e5)

# 获取总和
colSums(chunk_col_sums)
R 的可扩展数据处理

并行化 chunk.apply()

# 使用 chunk.apply 从 foo.csv 获取按行分块
chunk_col_sums <- chunk.apply("foo.csv",

 # 处理每个分块的函数
 function(chunk) {

   # 将分块转为数据框
   d <- dstrsplit(chunk, col_types = rep("numeric", 3), sep = ",")
   colSums(d)
 }, 
 # 使用 2 个处理器读取并处理数据
 CH.PARALLEL = 2)

# 获取总和
colSums(chunk_col_sums)
R 的可扩展数据处理

关于并行化的说明

  • 增加处理器数量不一定会加速代码
  • 在单机上增加处理器通常会出现收益递减
R 的可扩展数据处理

让我们来练习!

R 的可扩展数据处理

Preparing Video For Download...