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

