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

PySpark环境下高效读取批量多行压缩sdf化学分子文件的方法咨询

PySpark批量读取gzip压缩SDF文件的高效方案

核心思路:利用SDF的$$$$条目分隔符做行分组,结合成熟的化学信息学库解析单条分子,避免手动解析的冗余和错误,同时适配Spark分布式调度特性。

1 第一步:批量读取原始文件,聚合得到单分子完整文本

SDF格式以$$$$作为单条分子记录的结束标记,你可以根据你的文件大小选择对应的读取方式:

场景1:大量小SDF文件(单个文件<128M)

直接用wholeTextFiles读取每个文件完整内容后拆分,性能最高:

from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("SDFReader").getOrCreate()

# 直接读取所有gzip压缩SDF文件的完整内容,Spark自动解压
whole_files_rdd = spark.sparkContext.wholeTextFiles("/path/to/your/sdfs/*.sdf.gz")
# 按$$$$拆分得到每条分子的完整文本,保留所属文件路径和分组序号
mol_raw_df = whole_files_rdd.flatMap(
    lambda file_item: [
        (file_item[0], idx, mol_text.strip()) 
        for idx, mol_text in enumerate(file_item[1].split("$$$$")) 
        if mol_text.strip()
    ]
).toDF(["file_path", "mol_group", "mol_text"])

场景2:大体积SDF文件

如果单文件体积较大,按行读取后分组聚合更稳妥:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, collect_list, concat_ws, input_file_name, monotonically_increasing_id, sum, when
from pyspark.sql.window import Window

spark = SparkSession.builder.appName("SDFReader").getOrCreate()

raw_lines_df = spark.read.text("/path/to/your/sdfs/*.sdf.gz") \
    .withColumn("file_path", input_file_name()) \
    .withColumn("row_id", monotonically_increasing_id())

# 给分子打分组标记,每遇到$$$$就给分组序号+1
window_spec = Window.orderBy("row_id")
raw_lines_df = raw_lines_df.withColumn(
    "mol_group",
    sum(when(col("value") == "$$$$", 1).otherwise(0)).over(window_spec)
)

# 聚合得到单分子完整文本
mol_raw_df = raw_lines_df.groupBy("file_path", "mol_group") \
    .agg(concat_ws("\n", collect_list("value")).alias("mol_text"))

2 第二步:用成熟库解析单分子内容

不建议手动实现字段解析逻辑,SDF格式存在大量边界规则,用RDKit等成熟化学信息学库的解析能力,性能和准确率都更高,示例UDF如下:

from pyspark.sql.types import StructType, StructField, StringType, DoubleType
from rdkit import Chem
import io

# 定义解析后的返回结构,按需扩展你需要的字段
schema = StructType([
    StructField("mol_id", StringType(), True),
    StructField("shape_tanimoto", DoubleType(), True),
    StructField("color_tanimoto", DoubleType(), True),
    StructField("tanimoto_combo", DoubleType(), True),
    StructField("smiles", StringType(), True)
])

def parse_sdf_mol(mol_text):
    if not mol_text:
        return (None, None, None, None, None)
    try:
        # 用io.StringIO包装文本给SDMolSupplier读取
        suppl = Chem.SDMolSupplier()
        suppl.SetData(mol_text)
        mol = next(suppl, None)
        if not mol:
            return (None, None, None, None, None)
        # 提取属性和分子信息
        mol_id = mol.GetProp("_Name") if mol.HasProp("_Name") else None
        shape_t = float(mol.GetProp("ShapeTanimoto")) if mol.HasProp("ShapeTanimoto") else None
        color_t = float(mol.GetProp("ColorTanimoto")) if mol.HasProp("ColorTanimoto") else None
        combo_t = float(mol.GetProp("TanimotoCombo")) if mol.HasProp("TanimotoCombo") else None
        smiles = Chem.MolToSmiles(mol)
        return (mol_id, shape_t, color_t, combo_t, smiles)
    except Exception as e:
        return (None, None, None, None, None)

# 注册UDF
from pyspark.sql.functions import udf
parse_sdf_udf = udf(parse_sdf_mol, schema)

# 解析得到最终结构化数据
parsed_mol_df = mol_raw_df.withColumn("parsed", parse_sdf_udf(col("mol_text"))) \
    .select("file_path", "mol_group", "parsed.*")

3 性能优化建议

  • 依赖部署:确保所有Spark executor节点都安装了对应版本的RDKit,推荐用conda打包依赖后通过--archives参数提交作业,避免环境不一致问题。
  • 大文件处理:gzip本身是不可切分压缩格式,单个超大gzip SDF文件会导致单task处理瓶颈,建议先拆分为多个128M左右的小gzip SDF文件再提交处理,可充分利用集群资源。
  • 脏数据过滤:可新增过滤逻辑将解析失败的记录单独落盘排查,避免作业中断。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 01:06:02