如何建立ADLS与Kafka连接以在Databricks/ADF中处理流数据
ADLS流数据对接Kafka(Databricks/ADF平台)落地方案
前置准备
- 权限打通:给Databricks、ADF的服务主体分配目标ADLS路径的存储Blob数据参与者、存储队列数据参与者权限(流监听目录变更需要队列权限);提前在Kafka集群创建目标Topic,给服务主体开放Topic的生产权限,配置好SASL/SSL认证凭据,同VNet环境提前做VNet对等,跨网络环境放通对应端口白名单
- 规则梳理:提前确认ADLS路径下的文件格式(JSON/Parquet/CSV等)、分区规则、写入完成标记(推荐业务源写完文件后生成_SUCCESS标记,避免读取半写入的脏数据),统计单文件大小、数据生成吞吐量,提前规划Kafka分区数匹配写入压力
实现路径1:Databricks Structured Streaming 链路(适合高吞吐、低延迟、有复杂转换需求的场景)
该链路用Databricks原生Auto Loader对接ADLS增量文件,直接通过Structured Streaming写入Kafka,不需要额外中间组件,端到端延迟可以做到秒级,支持Exactly Once语义。
- 配置ADLS访问认证,直接通过ABFSS协议访问存储,不需要额外挂载存储,核心代码如下:
# 替换尖括号内的参数为实际环境值 spark.conf.set("fs.azure.account.auth.type.<存储账号名>.dfs.core.windows.net", "OAuth") spark.conf.set("fs.azure.account.oauth.provider.type.<存储账号名>.dfs.core.windows.net", "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider") spark.conf.set("fs.azure.account.oauth2.client.id.<存储账号名>.dfs.core.windows.net", "<服务主体客户端ID>") spark.conf.set("fs.azure.account.oauth2.client.secret.<存储账号名>.dfs.core.windows.net", "<服务主体密钥>") spark.conf.set("fs.azure.account.oauth2.client.endpoint.<存储账号名>.dfs.core.windows.net", "https://login.microsoftonline.com/<租户ID>/oauth2/token") # 用Auto Loader监听ADLS指定路径的新增文件 adls_source = spark.readStream \ .format("cloudFiles") \ .option("cloudFiles.format", "json") # 替换为实际文件格式,支持parquet/csv/avro等 .option("cloudFiles.useNotifications", "true") # 基于存储队列监听文件变更,比定时扫描延迟低、性能好 .option("cloudFiles.includeExistingFiles", "false") # 作业启动后仅处理新增文件,存量回溯可手动改为true .option("cloudFiles.cleanupMode", "MOVE_TO_TRASH") # 处理完成的文件标记自动清理,避免重复扫描 .option("badRecordsPath", "abfss://<日志容器>@<存储账号名>.dfs.core.windows.net/bad_records") # 解析失败的脏数据单独存储 .load("abfss://<业务容器>@<存储账号名>.dfs.core.windows.net/<流数据指定路径>")
- 按业务需求做流数据处理:可以直接做字段过滤、格式转换、维度维表关联,注意控制单批次处理逻辑耗时,不要在流作业里跑全量大表关联:
# 示例逻辑:过滤测试数据、提取核心业务字段 processed_df = adls_source.filter("is_test != true") \ .select("event_id", "user_id", "event_time", "event_type", "event_attrs")
- 写入目标Kafka Topic,注意Kafka要求消息值为二进制格式,需要先将结构化数据转为JSON序列化:
from pyspark.sql.functions import to_json, struct kafka_query = processed_df \ .select(to_json(struct(*processed_df.columns)).alias("value")) \ .writeStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "<Kafka broker地址列表>") \ .option("kafka.security.protocol", "SASL_SSL") \ .option("kafka.sasl.mechanism", "PLAIN") \ .option("kafka.sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username='$ConnectionString' password='<Kafka认证凭据>';") \ .option("kafka.acks", "all") # 保证消息写入不丢 .option("topic", "<目标Kafka Topic名>") \ .option("checkpointLocation", "abfss://<checkpoint专用容器>@<存储账号名>.dfs.core.windows.net/adls2kafka_cp") \ .trigger(processingTime="5 seconds") # 根据业务延迟要求调整触发间隔,最低可配置为百毫秒级 .start() kafka_query.awaitTermination()
注意:checkpoint路径为流作业存储偏移量、作业状态的专用路径,禁止多个流作业共用同一路径,非特殊情况不要手动删除路径下的文件,否则会导致重复消费或者漏消费
实现路径2:ADF低代码编排链路(适合逻辑简单、开发资源不足、延迟要求分钟级的场景)
该链路不需要写Spark代码,全可视化配置,运维门槛低,端到端延迟一般在10秒到分钟级。
- 配置存储事件触发器:在ADF里新建管道,绑定ADLS Blob创建事件触发器,监听指定路径下的文件生成事件,可配置过滤规则仅当_SUCCESS标记文件生成时再触发管道运行,避免读取半写入文件;同时配置防抖规则,比如设置每30秒或攒够50个文件触发一次管道,避免小文件过多导致管道频繁启动浪费资源
- 配置源数据集:源端选择ADLS Gen2,指定对应文件路径、格式,配置解析规则,比如跳过错误行、指定列的类型映射
- 配置转换逻辑(可选):简单的字段过滤、类型转换、值映射可以直接用ADF内置的映射数据流实现;如果有复杂处理逻辑,可以在管道里插入Databricks Notebook活动,调用提前写好的处理脚本做计算
- 配置Kafka接收器:接收器选择Kafka,填入Kafka集群地址、认证信息、目标Topic名称,配置并行写入数(和Kafka分区数保持一致即可,避免写入瓶颈),设置写入失败的重试次数、超时时间,失败数据自动落日志存储
上线运维要点
- 数据一致性校验:每天定时统计ADLS路径下的写入记录数、Kafka Topic的实际写入记录数,差值超过阈值时立刻排查作业状态、checkpoint是否损坏
- 监控告警:配置作业监控,Databricks流作业处理延迟超过1分钟、ADF管道失败率超过1%时触发告警
- 小文件优化:如果业务源写入ADLS的单文件大小普遍小于1MB,在处理环节加微批合并逻辑,攒一批小文件合并后再写入Kafka,减少Kafka的请求压力,避免集群阻塞
内容的提问来源于stack exchange,提问作者harshith
相关产品推荐
相关产品推荐

