如何在Databricks环境下合并包含不同sampleId的VCF文件
报错原因说明
标准VCF格式强制要求所有变异行的样本ID顺序完全统一,你当前的聚合逻辑只是把存在的样本基因型拼到数组里,不同行的genotypes数组包含的样本ID不一致,Glow写入时无法推断统一的样本列表,因此触发错误:
Task failed while writing rows
Cannot infer sample ids because they are not the same in every row.
解决方案
提供两种适配Databricks Glow环境的低开销方案,无需生成全量虚拟表:
方案1:写入标准兼容VCF,自定义缺失值填充
仅需一次轻量的样本ID提取操作,资源开销极低,输出文件符合通用VCF规范。
- 提取所有VCF文件的全局样本ID列表
import pyspark.sql.functions as F all_sample_ids = set() for file in file_list: # 仅提取样本ID,不加载全量变异数据,开销极小 sample_ids = [row[0] for row in spark.read.format("vcf").load(file) .select(F.explode("genotypes.sampleId")).distinct().collect()] all_sample_ids.update(sample_ids) all_sample_ids = sorted(list(all_sample_ids))
- 编写UDF按全局样本ID补全基因型,缺失值可自定义
from pyspark.sql.types import ArrayType, StructType, StructField, StringType, BooleanType, IntegerType # 对齐Glow读取的基因型结构体定义 genotype_schema = ArrayType(StructType([ StructField("sampleId", StringType()), StructField("phased", BooleanType()), StructField("calls", ArrayType(IntegerType())) ])) @F.udf(returnType=genotype_schema) def fill_missing_genotypes(existing_genotypes): existing_map = {g["sampleId"]: g for g in existing_genotypes} filled = [] for sid in all_sample_ids: if sid in existing_map: filled.append(existing_map[sid]) else: # 自定义缺失值示例:无调用标识./.对应calls为[None, None],可按需修改为其他标识 filled.append({ "sampleId": sid, "phased": False, "calls": [None, None] }) return filled
- 写入VCF时显式指定样本列表
spark.read.format("vcf")\ .option("flattenInfoFields", True)\ .load(file_list)\ .groupBy('contigName', 'start', 'end', 'referenceAllele', 'alternateAlleles', 'qual', 'filters','splitFromMultiAllelic')\ .agg(F.flatten(F.collect_list('genotypes')).alias('genotypes'))\ .withColumn("genotypes", fill_missing_genotypes("genotypes"))\ .write.mode("overwrite").format("vcf")\ .option("sampleIds", ",".join(all_sample_ids))\ .save("my_output_destination")
方案2:继续使用bigVCF,自定义缺失填充值
bigVCF本身支持配置缺失基因型的显示值,无需修改聚合逻辑,只需添加写入参数即可:
spark.read.format("vcf")\ .option("flattenInfoFields", True)\ .load(file_list)\ .groupBy('contigName', 'start', 'end', 'referenceAllele', 'alternateAlleles', 'qual', 'filters','splitFromMultiAllelic')\ .agg(F.sort_array(F.flatten(F.collect_list('genotypes'))).alias('genotypes'))\ .write.mode("overwrite").format("bigvcf")\ # 自定义缺失调用标识,示例为无调用./.,可按需替换为其他值 .option("missingGenotypeCall", "./.")\ .save("my_output_destination")
内容的提问来源于stack exchange,提问作者pinegulf
相关产品推荐
相关产品推荐

