如何使用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
相关产品推荐
相关产品推荐

