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

