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

如何在Spark中从半结构化文本文件加载并指定Schema生成DataFrame?

解决半结构化文本转Spark DataFrame(带预设Schema)的问题

首先,我得先拆解你给出的示例文本的结构,这样才能精准提取字段并匹配你的预设Schema。咱们一步步来:

1. 先明确预设Schema(以Python为例)

假设你的预设Schema包含这些字段,我先定义出来(你可以根据实际需求调整):

from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType

# 预设的Schema,和你需求匹配
review_schema = StructType([
    StructField("review_title", StringType(), nullable=True),
    StructField("username", StringType(), nullable=True),
    StructField("country", StringType(), nullable=True),
    StructField("review_date", DateType(), nullable=True),
    StructField("overall_rating", IntegerType(), nullable=True),
    StructField("review_content", StringType(), nullable=True),
    StructField("traveller_type", StringType(), nullable=True),
    StructField("cabin_flown", StringType(), nullable=True),
    StructField("route", StringType(), nullable=True),
    StructField("date_flown", DateType(), nullable=True),
    StructField("seat_comfort", IntegerType(), nullable=True),
    StructField("cabin_staff_service", IntegerType(), nullable=True),
    StructField("ground_service", IntegerType(), nullable=True)
])

2. 核心:用正则表达式提取半结构化字段

你的文本是典型的"键值对+自由文本"混合结构,最靠谱的方式是用正则匹配每行的各个字段。针对你的示例文本,我写了一个适配的正则,还加了字段提取的逻辑:

方式一:用RDD处理(适合复杂转换场景)

from datetime import datetime

def parse_review(line):
    import re
    # 匹配你文本结构的正则表达式,可根据实际文本调整
    pattern = r'^"([^"]+)" (\w+ \w+) \(([^)]+)\) (\w+ \w+ \d{4}) (\d+) ([\w\s\.]+?) Type Of Traveller ([\w\s]+) Cabin Flown (\w+) Route ([\w\s]+) Date Flown (\w+ \d{4}) Seat Comfort (\d) Cabin Staff Service (\d) Ground Service (\d) Value For .*$'
    match = re.match(pattern, line)
    
    if not match:
        # 处理不匹配的行,返回全None方便后续过滤
        return tuple([None]*13)
    
    # 提取各个字段
    review_title = match.group(1)
    username = match.group(2)
    country = match.group(3)
    
    # 转换评论日期格式
    review_date = datetime.strptime(match.group(4), "%dth %B %Y").date()
    overall_rating = int(match.group(5))
    review_content = match.group(6).strip()
    traveller_type = match.group(7)
    cabin_flown = match.group(8)
    route = match.group(9)
    
    # 转换飞行日期格式
    date_flown = datetime.strptime(match.group(10), "%B %Y").date()
    seat_comfort = int(match.group(11))
    cabin_staff_service = int(match.group(12))
    ground_service = int(match.group(13))
    
    return (review_title, username, country, review_date, overall_rating, review_content, traveller_type, cabin_flown, route, date_flown, seat_comfort, cabin_staff_service, ground_service)

# 读取文本文件为RDD
text_rdd = spark.sparkContext.textFile("/path/to/your/reviews.txt")
# 解析并过滤无效行
parsed_rdd = text_rdd.map(parse_review).filter(lambda x: all(field is not None for field in x))
# 转换为DataFrame并应用预设Schema
review_df = spark.createDataFrame(parsed_rdd, schema=review_schema)

方式二:用DataFrame内置函数(更高效,适合大数据场景)

如果你的数据量很大,推荐用Spark的内置函数直接在DataFrame层面处理,避免RDD的序列化开销:

from pyspark.sql.functions import regexp_extract, to_date

# 读取原始文本为DataFrame
raw_df = spark.read.text("/path/to/your/reviews.txt")

# 用regexp_extract提取每个字段,同时转换数据类型
review_df = raw_df.select(
    regexp_extract("value", r'^"([^"]+)"', 1).alias("review_title"),
    regexp_extract("value", r'" (\w+ \w+) \(', 1).alias("username"),
    regexp_extract("value", r'\(([^)]+)\)', 1).alias("country"),
    to_date(regexp_extract("value", r'\) (\w+ \w+ \d{4})', 1), "dth MMMM yyyy").alias("review_date"),
    regexp_extract("value", r'\d{4} (\d+)', 1).cast(IntegerType()).alias("overall_rating"),
    regexp_extract("value", r'(\d+) ([\w\s\.]+?) Type Of Traveller', 2).alias("review_content"),
    regexp_extract("value", r'Type Of Traveller ([\w\s]+) Cabin Flown', 1).alias("traveller_type"),
    regexp_extract("value", r'Cabin Flown (\w+) Route', 1).alias("cabin_flown"),
    regexp_extract("value", r'Route ([\w\s]+) Date Flown', 1).alias("route"),
    to_date(regexp_extract("value", r'Date Flown (\w+ \d{4})', 1), "MMMM yyyy").alias("date_flown"),
    regexp_extract("value", r'Seat Comfort (\d)', 1).cast(IntegerType()).alias("seat_comfort"),
    regexp_extract("value", r'Cabin Staff Service (\d)', 1).cast(IntegerType()).alias("cabin_staff_service"),
    regexp_extract("value", r'Ground Service (\d)', 1).cast(IntegerType()).alias("ground_service")
).filter(review_df.review_title.isNotNull())

3. 注意事项

  • 如果你的文本有格式变化(比如用户名带中间名、日期格式不同、评分是完整的1-5星级字符串),需要调整正则表达式的匹配规则,建议先抽样查看不匹配的行,针对性优化。
  • 日期转换时,要确保strptime的格式符和文本中的日期格式完全匹配,比如如果有"17 Oct 2015"这种缩写,要改成"%d %b %Y"。
  • 如果有大量无效行,可以先统计不匹配的比例,再决定是过滤还是补充正则规则。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:22:18