You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Scala在Databricks中实现CSV文件内嵌套数组对应字符串的查找替换

Databricks Scala 加载CSV并处理嵌套数组字符串查找替换实现方案

默认适配Databricks Runtime 7.x及以上版本,内置Spark SQL相关函数无需额外引入依赖。

1. 步骤1:加载CSV文件

// 读取CSV文件,可根据实际情况调整header、delimiter、quote等参数
val df = spark.read
  .option("header", "true")
  .option("inferSchema", "true")
  .option("delimiter", ",")
  .csv("dbfs:/path/to/your/file.csv") // 替换为你的CSV文件在DBFS/ADLS/S3的实际路径

// 查看原始数据结构与内容,确认嵌套数组字段的存储格式
df.printSchema()
df.show(false)

2. 步骤2:嵌套数组字符串查找替换实现

如果读取CSV后对应字段已经是原生Array类型(提前手动指定Schema的场景下),可直接跳过解析步骤,直接调用transform函数处理数组元素即可。

场景A:数组以JSON字符串格式存储在CSV中(最常见)

比如字段array_col的值为["old_str1","old_str2","other_str"]这类序列化JSON数组:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

// 定义数组的Schema,可根据实际元素类型调整
val arraySchema = ArrayType(StringType)

// 将JSON字符串解析为Spark原生Array类型
val parsedDf = df.withColumn("parsed_array", from_json(col("array_col"), arraySchema))

// 遍历数组元素执行查找替换,可自定义多个替换规则
val replacedDf = parsedDf.withColumn("replaced_array", 
  transform(col("parsed_array"), element => 
    when(element === "old_str1", "new_str1")
    .when(element === "old_str2", "new_str2")
    // 正则匹配替换示例:when(element.rlike("^old_prefix.*"), "new_value")
    .otherwise(element)
  )
)

// 处理完成后将数组转回CSV可存储的字符串格式,保留原字段名
val finalDf = replacedDf.withColumn("array_col", to_json(col("replaced_array")))
  .drop("parsed_array", "replaced_array")

场景B:数组以指定分隔符拼接的字符串格式存储在CSV中

比如字段array_col的值为old_str1|old_str2|other_str,用|作为分隔符:

import org.apache.spark.sql.functions._

// 按分隔符拆分字符串为数组
val splitDf = df.withColumn("split_array", split(col("array_col"), "\\|")) // 特殊分隔符需要转义

// 遍历数组执行替换
val replacedDf = splitDf.withColumn("replaced_array",
  transform(col("split_array"), element =>
    when(element === "old_str1", "new_str1")
    .when(element === "old_str2", "new_str2")
    .otherwise(element)
  )
)

// 拼接回原分隔符格式的字符串,保留原字段名
val finalDf = replacedDf.withColumn("array_col", array_join(col("replaced_array"), "|"))
  .drop("split_array", "replaced_array")

3. 结果导出(可选)

处理完成后可写回CSV存储:

finalDf.write
  .option("header", "true")
  .mode("overwrite") // 可根据需求替换为append/ignore等写入模式
  .csv("dbfs:/path/to/output/result.csv")

注意事项

  • 大规模数据场景下建议关闭inferSchema,手动指定Schema可大幅提升读取效率
  • 如果数组内是嵌套对象类型,只需调整from_json的Schema为对应的StructType即可适配
  • 字符串匹配可灵活使用contains、rlike等函数满足模糊替换需求

内容的提问来源于stack exchange,提问作者Ram

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 09:24:06