在R语言中更新Arrow数据集的最佳实践
更新Arrow数据集的最佳实践
问题场景
先创建按cyl分区的Arrow数据集:
suppressMessages(library(dplyr)) suppressMessages(library(arrow)) td <- tempdir() # 查看mtcars前6行数据 head(mtcars) #> mpg cyl disp hp drat wt qsec vs am gear carb #> Mazda RX4 21.0 6 160 110 3.90 2.620 16.46 0 1 4 4 #> Mazda RX4 Wag 21.0 6 160 110 3.90 2.875 17.02 0 1 4 4 #> Datsun 710 22.8 4 108 93 3.85 2.320 18.61 1 1 4 1 #> Hornet 4 Drive 21.4 6 258 110 3.08 3.215 19.44 1 0 3 1 #> Hornet Sportabout 18.7 8 360 175 3.15 3.440 17.02 0 0 3 2 #> Valiant 18.1 6 225 105 2.76 3.460 20.22 1 0 3 1
写入分区数据集:
write_dataset(mtcars, td, partitioning = "cyl")
此时生成3个分区文件夹:
dir(td) #> [1] "cyl=4" "cyl=6" "cyl=8"
尝试过滤cyl == 6的数据并重新写入原路径后,发现旧分区数据仍保留:
open_dataset(td) |> filter(cyl == 6) |> write_dataset(td, partitioning = "cyl") # 查看文件夹,旧分区仍存在 dir(td) #> [1] "cyl=4" "cyl=6" "cyl=8" # 数据集仍包含所有cyl值 open_dataset(td) |> distinct(cyl) |> collect() #> # A tibble: 3 × 1 #> cyl #> <int> #> 1 6 #> 2 4 #> 3 8
正确更新方式
1. 替换特定分区数据
如果仅需更新单个分区(如cyl=6),可先删除旧分区文件夹,再写入新数据:
# 提取需要更新的分区数据 updated_data <- open_dataset(td) |> filter(cyl == 6) |> collect() # 删除旧的cyl=6分区 unlink(file.path(td, "cyl=6"), recursive = TRUE) # 写入新的cyl=6数据到数据集路径 write_dataset(updated_data, td, partitioning = "cyl")
验证结果:
# 查看文件夹,仅cyl=4、cyl=6、cyl=8仍存在,但cyl=6已更新 dir(td) #> [1] "cyl=4" "cyl=6" "cyl=8" # 确认数据正确 open_dataset(td) |> filter(cyl == 6) |> collect()
2. 全量覆盖整个数据集
若要替换整个数据集(只保留新写入的分区),使用existing_data = "overwrite"参数:
open_dataset(td) |> filter(cyl == 6) |> write_dataset(td, partitioning = "cyl", existing_data = "overwrite")
验证结果:
# 仅保留cyl=6分区 dir(td) #> [1] "cyl=6" # 数据集仅包含cyl=6的数据 open_dataset(td) |> distinct(cyl) |> collect() #> # A tibble: 1 × 1 #> cyl #> <int> #> 1 6
3. 增量追加数据到分区
如果需要在现有分区中追加新数据(而非替换),使用默认的existing_data = "append"参数即可,确保新数据的分区键与目标分区匹配:
# 模拟新的cyl=6数据 new_data <- tibble(mpg = 25, cyl = 6, disp = 150, hp = 100, drat = 3.8, wt = 2.5, qsec = 16.5, vs = 1, am = 1, gear = 4, carb = 2) # 追加到现有数据集 write_dataset(new_data, td, partitioning = "cyl", existing_data = "append")
验证结果:
# cyl=6分区数据量增加 open_dataset(td) |> filter(cyl == 6) |> count() |> collect()
内容的提问来源于stack exchange,提问作者Philippe Massicotte
相关产品推荐
相关产品推荐

