在Foundry环境中拆分PySpark数据集TSV格式单列至多列
在Foundry中拆分TSV格式的PySpark数据集列
问题背景
我在Foundry环境中有一个PySpark数据集,仅包含名为data的列,该列每行都是制表符分隔(TSV)格式的字符串,示例如下:
ott-akamai-logs-processor srv 2023-07-29 17:46:50.134 2023-07-29 17:46:49.358 unstruct 103b9271-777 ott node-3.13.1 ssc-2.8.2-kinesis snowplow-enrich-kinesis-3.7.0 3.65.234.x 12345679 DE HE Karachi 60313 50.1188 8.6843 Malta {"schema":"iglu:com.xxx/1-0-0","data":{"schema":"xxx/hls_manifest_requested/jsonschema/1-0-1","data":{"channel":"bildtv-broadcast","session_id":"xxx","request_id":"xxx","total_bytes":351,"referrer":"^","geo_country":"DE","geo_state":"Berlin","geo_city":"-","variant_name":"6.m3u8"}}} snowplow-nodejs-tracker/3.13.1 Europe/Berlin 2023-07-29 17:46:49.281 {"schema":"xxx/contexts/jsonschema/1-0-1","data":[{"schema":"iglu:nl.basjes/yauaa_context/jsonschema/1-0-4","data":{"deviceBrand":"Unknown","deviceName":"Unknown","operatingSystemVersionMajor":"??","layoutEngineNameVersion":"Unknown ??","operatingSystemNameVersion":"Unknown ??","agentInformationEmail":"Unknown","networkType":"Unknown","webviewAppNameVersionMajor":"Unknown ??","layoutEngineNameVersionMajor":"Unknown ??","operatingSystemName":"Unknown","agentVersionMajor":"3","layoutEngineVersionMajor":"??","webviewAppName":"Unknown","deviceClass":"Unknown","agentNameVersionMajor":"Snowplow-Nodejs-Tracker 3","operatingSystemNameVersionMajor":"Unknown ??","webviewAppVersionMajor":"??","operatingSystemClass":"Unknown","webviewAppVersion":"??","layoutEngineName":"Unknown","agentName":"Snowplow-Nodejs-Tracker","agentVersion":"3.13.1","layoutEngineClass":"Unknown","agentNameVersion":"Snowplow-Nodejs-Tracker 3.13.1","operatingSystemVersion":"??","agentClass":"Special","layoutEngineVersion":"??","agentInformationUrl":"Unknown"}},{"schema":"iglu:com.snowplowanalytics.snowplow/ua_parser_context/jsonschema/1-0-0","data":{"useragentFamily":"Other","useragentMajor":null,"useragentMinor":null,"useragentPatch":null,"useragentVersion":"Other","osFamily":"Other","osMajor":null,"osMinor":null,"osPatch":null,"osPatchMinor":null,"osVersion":"Other","deviceFamily":"Other"}}]} 2023-07-29 17:46:09.938 com.axelspringer.ott hls_manifest_requested jsonschema 1-0-1 2023-07-29 17:46:09.938
需要在给定的函数框架内,将这些制表符分隔的内容拆分到不同列:
def unnamed_1(my_df): df = my_df return df
解决方案
步骤说明
- 拆分TSV列:用PySpark的
split函数,以制表符\t为分隔符拆分data列,得到数组类型的列。 - 展开数组为多列:通过
select将数组中的每个元素映射为单独列,并为每列指定有意义的名称(需匹配TSV字段顺序)。 - 处理连续制表符:示例存在连续制表符,拆分后会产生空字符串,可按需保留或清理这些空值。
完整代码实现
from pyspark.sql import functions as F def unnamed_1(my_df): # 定义TSV对应的列名(需根据实际字段顺序调整,示例为参考) column_names = [ "processor_name", "service_type", "event_time", "collector_time", "event_type", "event_id", "platform", "tracker_version", "ssc_version", "enrich_version", "ip_address", "user_id", "country_code", "state_code", "city", "postal_code", "lat", "lon", "region", "unstruct_event", "user_agent", "timezone", "dvce_created_time", "contexts", "true_tstamp", "vendor", "event_name", "schema_type", "schema_version", "etl_tstamp" ] # 拆分data列为数组 split_col = F.split(my_df["data"], "\t") # 展开数组为多个列,并指定列名 df = my_df.select( *[split_col.getItem(i).alias(col_name) for i, col_name in enumerate(column_names)] ) # 可选:过滤所有字段为空的行(按需调整) # df = df.filter(F.length(F.trim(F.concat_ws("", *df.columns))) > 0) return df
关键说明
- 列名匹配:
column_names需严格对应TSV字符串的字段顺序,可通过分析样本数据的字段数量和含义调整。 - 空值处理:
split默认保留连续分隔符产生的空字符串,若不需要可后续用filter或withColumn清理。 - 类型转换:拆分后的列均为字符串类型,如需转换时间、经纬度等字段类型,可添加
cast操作,例如:df = df.withColumn("event_time", F.to_timestamp(F.col("event_time"))) df = df.withColumn("lat", F.col("lat").cast("double"))
内容的提问来源于stack exchange,提问作者x89
相关产品推荐
相关产品推荐

