如何使用PySpark读取无换行管道分隔文本生成指定结构DataFrame
PySpark 实现单行多管道分隔文本转结构化DataFrame
实现步骤
- 先读取原始TXT文件,获取整行的原始文本内容
- 对文本做分割清洗后按固定长度分组,匹配每行3个字段的结构
- 基于分组后的结构化数据直接构建DataFrame
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 初始化SparkSession spark = SparkSession.builder.appName("SplitPipeTxt").getOrCreate() # 1. 读取原始txt文件,读取为单行字符串RDD,替换为你的实际文件路径 raw_rdd = spark.sparkContext.textFile("./test_data.txt") # 2. 处理文本:按|分割、过滤空值、每3个字段为一组拆分 def process_line(line): # 按|分割后过滤空字符串,得到所有字段列表 fields = [f.strip() for f in line.split("|") if f.strip()] # 每3个元素切分为一组,对应一行的c0、c1、c2 return [fields[i:i+3] for i in range(0, len(fields), 3)] # 扁平化处理得到每行的字段列表 row_rdd = raw_rdd.flatMap(process_line) # 3. 定义表结构Schema schema = StructType([ StructField("c0", IntegerType(), True), StructField("c1", StringType(), True), StructField("c2", IntegerType(), True) ]) # 4. 构建DataFrame df = spark.createDataFrame(row_rdd, schema=schema) # 验证输出 df.show()
按「替换第4个|为换行符」思路的实现方案
如果你需要严格按照你提到的思路实现,可参考如下处理逻辑:
def replace_4th_pipe(line): fields = line.split("|") new_fields = [] count = 0 for f in fields: count +=1 # 每遇到第4个|就替换为换行符 if count %4 ==0: new_fields.append(f + "\n") else: new_fields.append(f + "|") # 拼接后按换行拆分,过滤首尾空行 processed_lines = "".join(new_fields).split("\n") return [l.strip("|") for l in processed_lines if l.strip()] # 处理后按|拆分每行得到字段列表 processed_rdd = raw_rdd.flatMap(replace_4th_pipe).map(lambda x: [i.strip() for i in x.split("|")]) df = spark.createDataFrame(processed_rdd, schema=schema)
输出验证
两种方案运行后得到的DataFrame均符合要求:
+---+---------+---+ | c0| c1| c2| +---+---------+---+ | 1| NewYork| 96| | 2| Chennai| 84| | 3|Amsterdam| 96| +---+---------+---+
内容的提问来源于stack exchange,提问作者Karthik Deepan
相关产品推荐
相关产品推荐

