如何在Palantir Foundry中为解析后的数据集添加文件名列?
Palantir Foundry中解析多CSV并添加对应文件名列的解决方案
问题描述
在Palantir Foundry中拥有包含多个CSV文件的原始数据集,需要实现两个目标:
- 将CSV文件解析为数据集
- 添加包含对应文件名的新列
熟悉PySpark但对平台不熟悉,当前代码能完成解析,但所有行的文件名都显示为同一个,无法匹配各自的源文件。
问题根源
原代码中通过list(raw.filesystem().ls(glob='*.csv'))[0].path仅获取了第一个CSV文件的路径,再用F.lit(file_name)将这个固定值赋值给所有行,导致所有行的文件名完全一致。
解决方案
利用Spark内置的input_file_name()函数,该函数可以自动为每行数据匹配其对应的源文件完整路径,完美解决多文件的文件名关联问题。修改后的代码如下:
from transforms.api import transform, Input, Output, incremental from transforms.verbs.dataframes import sanitize_schema_for_parquet from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType @incremental() @transform( output=Output("rid"), raw=Input("rid") ) def read_csv(ctx, raw, output): filesystem = raw.filesystem() hadoop_path = filesystem.hadoop_path files = [f"{hadoop_path}/{f.path}" for f in filesystem.ls()] csv_schema = StructType([ StructField("SamaccountName", StringType(), True), StructField("DisplayName", StringType(), True), StructField("Alias", StringType(), True), StructField("PrimarysmtpAddress", StringType(), True), StructField("TotalMBXSize", StringType(), True), StructField("UserMailboxSize", StringType(), True), StructField("TotalDeletedItemSize", StringType(), True), StructField("ItemCount", IntegerType(), True), StructField("LocatedOnDatabase", StringType(), True), StructField("CreationDateTime", TimestampType(), True), StructField("MailboxType", StringType(), True), StructField("RecipientTypeDetails", StringType(), True), StructField("RetentionPolicy", StringType(), True), StructField("CustomAttribute3", StringType(), True), StructField("CustomAttribute7", StringType(), True) ]) df = ( ctx .spark_session .read # .option("encoding", "UTF-8") .option("header", True) .schema(csv_schema) .option("delimiter", ";") .csv(files) ).withColumn("full_file_path", F.input_file_name()) # 获取每行对应的完整文件路径 # 从完整路径中提取纯文件名(如果不需要完整路径可保留此步骤,否则直接用full_file_path即可) df = df.withColumn("file_name", F.element_at(F.split(df.full_file_path, "/"), -1)) df = df.withColumn( "Date", F.to_date( F.regexp_extract(df.file_name, r'Outlookstats_(\d{8})_\d{6}\.csv', 1), 'yyyyMMdd')) output.write_dataframe(sanitize_schema_for_parquet(df))
关键修改说明
- 移除了固定获取单个文件名的逻辑,替换为
F.input_file_name(),自动关联每行数据的源文件路径 - 通过
F.split和F.element_at从完整路径中提取纯文件名,若需要完整路径可跳过此步骤 - 保留原有日期提取逻辑,现在
file_name列已正确对应每行的源文件,日期提取结果也会准确匹配
额外提示
- 因使用
incremental()装饰器,增量运行时会自动处理新增的CSV文件,无需额外调整 - 原代码同时指定了
schema(csv_schema)和option("inferSchema", True),可移除inferSchema,显式指定的Schema优先级更高,能避免不必要的类型推断开销
内容的提问来源于stack exchange,提问作者Pankaj Nain
相关产品推荐
相关产品推荐

