Apache Spark提取HDFS中CSV列冒号前子串的实现方法求助
提取CSV中title列冒号前内容的Spark解决方案
嘿,这个需求我熟!你要从类似EMS: BACK PAINS/INJURY的title列里只提取冒号前的部分(比如EMS),完全可以用Spark自带的字符串处理函数轻松实现,不用绕弯路。下面给你两种最常用的方法,附完整代码示例:
方法1:用substring_index(推荐,最直接)
这个函数就是专门用来按分隔符截取子串的,语法简单,处理效率也高。我们只需要指定分隔符为:,并取第一个分隔符前的内容就行:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.{StructType, StructField, StringType} object EmergencyTypeExtractor { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("EmergencyTypeExtractor") .master("local[*]") // 本地测试用,生产环境记得移除 .getOrCreate() // 定义你的schema(这里假设包含title和其他列,按需调整) val schema = StructType(Array( StructField("title", StringType, nullable = true), StructField("event_time", StringType, nullable = true), // 其他列... StructField("location", StringType, nullable = true) )) // 从HDFS加载CSV到DataFrame val rawDf = spark.read .schema(schema) .option("header", "true") // 如果CSV文件有表头就打开这个选项 .csv("hdfs://your/hdfs/path/to/data.csv") // 提取冒号前的内容,生成新列emergency_type;如果想替换原title列,直接把新列名改成title就行 val processedDf = rawDf.withColumn("emergency_type", substring_index(col("title"), ":", 1)) // 打印结果验证 processedDf.select("title", "emergency_type").show(false) spark.stop() } }
方法2:用split函数分割字符串
另一种思路是把title列按:分割成数组,然后取数组的第一个元素:
// 替换上面的processedDf行即可 val processedDf = rawDf.withColumn("emergency_type", split(col("title"), ":")(0))
额外处理:兼容特殊情况
如果你的数据里存在冒号前后有空格的情况(比如EMS : BACK PAINS),可以先做个正则替换清理空格,再提取:
val processedDf = rawDf.withColumn( "emergency_type", substring_index(regexp_replace(col("title"), ":\\s*", ":"), ":", 1) )
这样不管冒号后面有没有空格,都能准确提取到冒号前的内容啦。如果title列里有不含冒号的行,这两种方法都会直接返回原字符串,不用担心报错~
内容的提问来源于stack exchange,提问作者tejas zinzuwadiya
相关产品推荐
相关产品推荐

