如何通过Azure Data Factory调用Spark实现两个JSONL文件的加载转换合并
Azure Data Factory结合Spark实现双JSONL文件关联转换操作指南
前提准备
- 已创建Azure Data Factory实例,且拥有创建管道、链接服务的对应权限
- 已创建可对接ADF的Spark计算资源:推荐使用Azure Synapse Spark池,或是Azure Databricks工作区,两种资源均支持对接ADF的Spark活动
- 两个待处理JSONL文件已上传到Azure存储账户(Blob存储/ADLS Gen2均可),提前记录两个文件的存储路径,以及要用于关联的共有字段名
- 已在ADF中创建存储账户的链接服务,保障Spark活动可以正常访问存储中的文件
步骤1:在ADF管道中创建Spark活动
- 进入ADF Studio的「作者」面板,新建空白管道
- 在活动面板的「计算」分类下,选择和你所用Spark资源对应的活动:使用Synapse Spark池选「Synapse Spark 活动」,使用Databricks选「Databricks Spark 活动」,拖拽到管道画布上
- 填写活动基本名称后,在「Azure Synapse」/「Databricks」标签页选择你提前创建好的Spark计算资源链接服务完成绑定
步骤2:编写Spark处理代码(核心部分)
Spark活动的代码可以直接内嵌在ADF活动中,也可以上传到存储账户后指定路径,以下是可直接复用的代码示例,替换对应占位符即可使用:
# 导入Spark会话 from pyspark.sql import SparkSession # 初始化Spark会话 spark = SparkSession.builder.appName("JSONLJoinProcess").getOrCreate() # ------------------- 替换以下参数为你的实际信息 ------------------- # 第一个JSONL文件的存储路径,ADLS Gen2用abfss协议,Blob存储用wasb协议 file1_path = "abfss://你的容器名@你的存储账户名.dfs.core.windows.net/第一个文件路径/xxx.jsonl" # 第二个JSONL文件的存储路径 file2_path = "abfss://你的容器名@你的存储账户名.dfs.core.windows.net/第二个文件路径/yyy.jsonl" # 两个文件共有的关联字段名 join_column = "你要用于关联的共有字段名" # 合并后文件的输出存储路径 output_path = "abfss://你的容器名@你的存储账户名.dfs.core.windows.net/输出文件路径/" # ----------------------------------------------------------------- # 加载两个JSONL文件,Spark原生支持JSONL格式,直接读取即可 df1 = spark.read.json(file1_path, multiLine=False) df2 = spark.read.json(file2_path, multiLine=False) # 执行JOIN操作,默认是内关联,需要左/右/全关联可将how参数修改为'left'/'right'/'full' joined_df = df1.join(df2, on=join_column, how='inner') # 输出合并后的数据为JSONL格式,mode参数设为'overwrite'表示覆盖已有文件,'append'表示追加 joined_df.write.mode("overwrite").json(output_path)
注意:如果你的JSONL文件存在特殊编码、字段嵌套等情况,可在
read.json方法中补充对应参数即可,比如添加encoding='utf-8'指定文件编码
步骤3:配置Spark活动的代码参数
- 若使用Synapse Spark活动:在活动的「设置」标签页,可选择「文件路径」提前将代码上传到存储中调用,也可以选择「内联」直接将上述代码粘贴到输入框中
- 若使用Databricks Spark活动:在活动的「设置」标签页选择「Python 任务」,指定代码文件路径或者填写内联代码即可
步骤4:调试运行与结果验证
- 点击管道顶部的「调试」按钮,触发管道运行
- 运行完成后到你指定的
output_path路径下查看生成的文件,即为JOIN完成后的合并数据 - 若运行报错,可直接查看Spark活动的输出日志,可快速定位是路径权限问题、字段名不匹配问题还是代码语法问题
内容的提问来源于stack exchange,提问作者Bruno
相关产品推荐
相关产品推荐

