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

在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

解决方案

步骤说明

  1. 拆分TSV列:用PySpark的split函数,以制表符\t为分隔符拆分data列,得到数组类型的列。
  2. 展开数组为多列:通过select将数组中的每个元素映射为单独列,并为每列指定有意义的名称(需匹配TSV字段顺序)。
  3. 处理连续制表符:示例存在连续制表符,拆分后会产生空字符串,可按需保留或清理这些空值。

完整代码实现

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 22:25:54