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

如何在Spark中一次性输出JSON数组字典?Cosmos配置报错排查

问题解决:PySpark输出JSON数组格式 + Cosmos配置错误排查

一、修改代码实现JSON数组格式输出

原代码逐条处理并记录单条键值对,需改为收集所有符合条件的记录后统一输出目标JSON结构,修改后的代码如下:

from pyspark.sql import SparkSession
import json
from pyspark.sql.functions import explode

class ReadParams:
    def __init__(self):
        self.spark_session = (
            SparkSession.builder
            .appName("read parameters")
            .getOrCreate()
        )
        self.input_directory = dbutils.secrets.get(scope="test", key="inpath")
        # 初始化列表用于收集结果
        self.type_list = []
        self.container_list = []

    def readConfig(self):
        logger.info(f"Reading the Config present in {self.input_directory} ")
        # 匹配输入示例的根节点"Input",修正字段映射
        dfConfig = self.spark_session.read.option("multiline","true") \
                       .json(self.input_directory)
        dfConfig = dfConfig.select(explode('Input').alias('input'))\
             .select('input.Type','input.container','input.Query')
        
        for f in dfConfig.rdd.toLocalIterator():
            self.Type = f[0]
            self.container = f[1]
            self.query = f[2]
            self.readFromCosmos()
        
        # 所有处理完成后生成目标JSON并输出
        result = {
            "Type": self.type_list,
            "container": self.container_list
        }
        print(json.dumps(result, indent=2))

    def readFromCosmos(self):
        Endpoint = dbutils.secrets.get(scope="test", key="Endpoint")
        MasterKey = dbutils.secrets.get(scope="test", key="DevKey")
        cosmosDatabaseName = "data_v2"
        cosmosContainerName = self.container
        query = self.query
        cfg = {
            "spark.cosmos.accountEndpoint": Endpoint,
            "spark.cosmos.accountKey": MasterKey,
            "spark.cosmos.database": cosmosDatabaseName,
            "spark.cosmos.read.inferSchema.enabled": "true",
            "spark.cosmos.container": cosmosContainerName,
            "spark.cosmos.read.customQuery": query,
        }
        
        logger.info(f"Captured Cosmos Configs {cfg}")
        # 移除重复配置,避免冗余
        df = self.spark_session.read.format("cosmos.oltp").options(**cfg).load()
        df.cache()
        
        if len(df.head(1)) > 0:
            self.Status = "Failure"
            # 将符合条件的记录加入列表
            self.type_list.append(self.Type)
            self.container_list.append(self.container)

if __name__ == "__main__":
   p = ReadParams()
   p.readConfig() 

关键修改说明:

  • 新增type_list和container_list统一收集符合条件的Type和container值
  • 修正输入JSON的根节点映射(原代码用a,输入示例为Input)及字段索引错误
  • 移除重复的inferSchema配置,简化代码逻辑
  • 所有数据处理完成后批量生成并打印目标JSON结构

运行后输出将符合期望格式:

{
  "Type": ["Completeness", "Accuracy"],
  "container": ["ct-data", "trans-data"]
}

二、Cosmos配置错误原因确认

错误IllegalArgumentException: The config property 'spark.cosmos.read.customquery' is invalid. No config setting with this name exists.确实由Cosmos Connector导致,核心原因是配置参数大小写不匹配:

  • Cosmos DB Spark Connector的配置参数区分大小写,正确参数名为spark.cosmos.read.customQuery(驼峰命名,Q大写)
  • 错误信息中参数为全小写的customquery,说明代码存在拼写错误,或有逻辑将参数名转为小写,导致Connector无法识别该配置

解决方法:

  • 严格按照官方规范使用驼峰命名书写所有Cosmos配置参数,如spark.cosmos.read.customQuery
  • 检查是否有参数传递、配置加载逻辑会自动转换参数大小写,修正该逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 15:25:34