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

如何通过ADF v2基于多层级Salesforce对象生成JSON并存储到Blob

多层级Salesforce数据生成嵌套JSON到Blob的实现方案

方案一:Azure Data Factory Data Flow 多层嵌套聚合(可视化配置)

适合中等数据量,无需代码,通过可视化转换逐步从底层往顶层聚合:

  1. 导入所有数据源:

    • 在Data Flow中创建5个Salesforce源,分别读取Program、Project、Items、Account、Contact对象,确保所有关联字段(如ProgramId、AccountId、ProjectId)都被选中。
  2. 聚合Items到Project:

    • 添加Aggregate转换,以ProjectId为分组键,使用collect()函数将Items字段聚合成数组:
      items: collect(@(itemId=ItemId, itemName=ItemName, ...))
      
    • 保留Project的其他字段(如ProgramId、AccountId、ProjectName),得到包含Items数组的Project数据集。
  3. 聚合Contact到Account:

    • 添加Aggregate转换,以AccountId为分组键,聚合Contact为数组:
      contacts: collect(@(contactId=ContactId, contactName=ContactName, ...))
      
    • 保留Account的其他字段(如ProgramId、AccountName),得到包含Contacts数组的Account数据集。
  4. 关联Project到Account并聚合:

    • 添加Join转换,将步骤2的Project数据集与步骤3的Account数据集通过AccountId关联。
    • 再添加Aggregate转换,以AccountId为分组键,将关联后的Project聚合成数组:
      projects: collect(@(projectId=ProjectId, projectName=ProjectName, items=items))
      
    • 保留Account的contacts数组和基础字段,得到完整的Account数据集(含Contacts、Projects-Items)。
  5. 关联Account和Project到Program并聚合:

    • 分别将步骤4的Account数据集、步骤2的Project数据集与Program数据集通过ProgramId做Join。
    • 添加Aggregate转换,以ProgramId为分组键,聚合Accounts和Projects数组:
      accounts: collect(@(accountId=AccountId, accountName=AccountName, contacts=contacts, projects=projects))
      projects: collect(@(projectId=ProjectId, projectName=ProjectName, items=items))
      
    • 按目标JSON结构调整聚合字段。
  6. 输出到Blob存储:

    • 添加Sink转换,选择Blob存储,格式设为JSON,根据需求选择输出为单文档或文档数组,配置存储路径和文件名规则。

方案二:Synapse Notebook(Spark脚本)处理

适合大数据量或复杂嵌套逻辑,灵活性更高:

  1. 读取Salesforce数据为DataFrame:

    # 初始化Spark读取Salesforce对象
    df_program = spark.read.format("salesforce").option("objectName", "Program").load()
    df_project = spark.read.format("salesforce").option("objectName", "Project").load()
    df_items = spark.read.format("salesforce").option("objectName", "Items").load()
    df_account = spark.read.format("salesforce").option("objectName", "Account").load()
    df_contact = spark.read.format("salesforce").option("objectName", "Contact").load()
    
  2. 逐步嵌套聚合:

    • 聚合Items到Project:
      from pyspark.sql.functions import collect_list, struct
      
      df_project_with_items = df_project.join(df_items, df_project.ProjectId == df_items.ProjectId, "left")\
          .groupBy(df_project.ProjectId, df_project.ProgramId, df_project.AccountId, df_project.ProjectName)\
          .agg(collect_list(struct(df_items.ItemId, df_items.ItemName)).alias("items"))
      
    • 聚合Contact到Account:
      df_account_with_contacts = df_account.join(df_contact, df_account.AccountId == df_contact.AccountId, "left")\
          .groupBy(df_account.AccountId, df_account.ProgramId, df_account.AccountName)\
          .agg(collect_list(struct(df_contact.ContactId, df_contact.ContactName)).alias("contacts"))
      
    • 关联Project到Account:
      df_account_full = df_account_with_contacts.join(df_project_with_items, df_account_with_contacts.AccountId == df_project_with_items.AccountId, "left")\
          .groupBy(df_account_with_contacts.AccountId, df_account_with_contacts.ProgramId, df_account_with_contacts.AccountName, df_account_with_contacts.contacts)\
          .agg(collect_list(struct(df_project_with_items.ProjectId, df_project_with_items.ProjectName, df_project_with_items.items)).alias("projects"))
      
    • 关联Account和Project到Program:
      df_program_full = df_program.join(df_account_full, df_program.ProgramId == df_account_full.ProgramId, "left")\
          .join(df_project_with_items, df_program.ProgramId == df_project_with_items.ProgramId, "left")\
          .groupBy(df_program.ProgramId, df_program.ProgramName)\
          .agg(collect_list(struct(df_account_full.AccountId, df_account_full.AccountName, df_account_full.contacts, df_account_full.projects)).alias("accounts"),
               collect_list(struct(df_project_with_items.ProjectId, df_project_with_items.ProjectName, df_project_with_items.items)).alias("projects"))
      
  3. 写入Blob存储:

    # 写入多文件JSON
    df_program_full.write.format("json").mode("overwrite").save("abfss://<container>@<storage-account>.dfs.core.windows.net/<path>")
    
    # 写入单文件JSON(适合小数据量)
    df_program_full.coalesce(1).write.format("json").mode("overwrite").save("abfss://<container>@<storage-account>.dfs.core.windows.net/<single-file-path>")
    

方案三:Salesforce SOQL嵌套查询+Copy Activity(轻量快捷)

适合小数据量,直接通过SOQL获取嵌套结构,无需复杂转换:

  1. 编写嵌套SOQL查询:

    SELECT ProgramId, ProgramName,
      (SELECT AccountId, AccountName,
        (SELECT ContactId, ContactName FROM Contacts),
        (SELECT ProjectId, ProjectName, (SELECT ItemId, ItemName FROM Items) FROM Projects)
      FROM Accounts),
      (SELECT ProjectId, ProjectName, (SELECT ItemId, ItemName FROM Items) FROM Projects)
    FROM Program
    

    注意:Salesforce对SOQL子查询有单批次记录数限制(默认最多5000条),数据量大时可能出现截断。

  2. Copy Activity导出:

    • 在ADF中创建Copy Activity,源选Salesforce,将上述SOQL作为自定义查询。
    • 目标选Blob存储,格式设为JSON,配置存储路径即可直接导出嵌套结构的JSON文件。

方案选择建议

  • 小数据量(单查询返回数<5000):优先用方案三,配置最快
  • 中等数据量:用方案一,可视化配置易维护
  • 大数据量/复杂逻辑:用方案二,Spark处理能力更强

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 20:07:11