如何通过Azure Data Factory连接本地Kafka并加载数据至ADLS
通过Azure Data Factory连接本地Kafka并将数据加载至ADLS的实现步骤
前提条件
- 本地Kafka集群正常运行,目标读取Topic已创建,Broker节点可被内网访问
- 已创建Azure Data Factory(ADF)实例
- 已创建Azure Data Lake Storage(ADLS)Gen2存储账户,且具备该账户的写入权限
- 部署**自托管集成运行时(Self-hosted Integration Runtime)**的本地服务器:需同时能访问本地Kafka集群和ADLS(ADLS可通过公网或Azure Private Link访问)
步骤1:配置自托管集成运行时(IR)
自托管IR是ADF访问本地私网资源的核心组件,配置步骤如下:
- 在ADF Studio进入「管理」→「集成运行时」,点击「新建」
- 选择「自托管」类型,完成名称和描述配置后点击「创建」
- 下载自托管IR安装包至本地服务器,运行安装程序并使用ADF生成的注册密钥完成注册
- 验证IR状态为「运行中」,确保该服务器能ping通Kafka Broker节点、访问Kafka端口(默认9092,SSL认证为9093)
步骤2:创建Kafka源数据集
- 在ADF Studio进入「数据」→「数据集」,点击「新建数据集」
- 搜索并选择「Kafka」连接器,点击「继续」
- 配置连接属性:
- 集成运行时:选择已部署的自托管IR
- Broker列表:填写本地Kafka Broker的内网地址(如
192.168.1.10:9092,192.168.1.11:9092) - Topic名称:选择要读取的目标Kafka Topic
- 认证类型:根据本地Kafka配置选择(无认证/SASL PLAIN/SASL SSL等),若需SSL需上传证书文件
- 配置格式设置:选择Kafka消息的格式(如JSON、Avro、CSV),并对应配置格式参数(比如JSON的编码、是否包含表头)
- 保存数据集,命名为
Kafka_Source
步骤3:创建ADLS目标数据集
- 在「数据集」中点击「新建数据集」,搜索并选择「Azure Data Lake Storage Gen2」,点击「继续」
- 选择数据存储类型(如「容器」),点击「继续」
- 配置连接属性:
- 集成运行时:若ADLS通过公网访问可选择托管IR,若用Private Link则选自托管IR
- 存储账户名称:选择目标ADLS Gen2账户
- 身份验证方式:选择「账户密钥」或「服务主体」(需确保服务主体拥有ADLS的「存储Blob数据贡献者」角色)
- 容器名称:选择目标容器,可指定文件夹路径
- 配置格式设置:选择与Kafka源匹配的输出格式(如Parquet、JSON),配置写入模式(追加/覆盖)
- 保存数据集,命名为
ADLS_Sink
步骤4:创建数据加载管道(两种方案)
方案1:使用复制活动(快速批量/增量加载)
- 进入「管道」→「新建管道」,添加「复制活动」至画布
- 配置复制活动的「源」选项卡:
- 数据集:选择
Kafka_Source - 读取选项:根据需求选择「从最新偏移量开始」「从最早偏移量开始」或「指定偏移量/时间戳」(增量加载推荐用偏移量或时间戳)
- 批量设置:可调整每次读取的消息数量(如
10000)优化性能
- 数据集:选择
- 配置复制活动的「接收器」选项卡:
- 数据集:选择
ADLS_Sink - 写入行为:选择「追加」(增量加载)或「覆盖」(全量加载)
- 文件命名:可设置前缀(如
kafka_data_)和滚动机制(按文件大小或行数)
- 数据集:选择
- 保存管道,点击「调试」测试运行,验证数据是否成功写入ADLS
方案2:使用数据流(需数据转换场景)
- 进入「数据流」→「新建数据流」,添加「源」转换
- 配置源转换:
- 数据源类型:选择「Kafka」
- 数据集:选择
Kafka_Source - 读取选项:同复制活动,支持增量读取配置
- 添加必要的转换组件(如「派生列」「筛选」「聚合」)处理数据(若无需转换可跳过)
- 添加「接收器」转换:
- 数据集:选择
ADLS_Sink - 写入模式:选择「追加」或「覆盖」
- 输出设置:配置文件分区、压缩格式等(如Parquet压缩选Snappy)
- 数据集:选择
- 返回管道界面,添加「执行数据流」活动,选择创建好的数据流
- 调试运行管道,验证转换后的数据是否写入ADLS
关键注意事项
- 网络连通性:确保自托管IR服务器与Kafka集群、ADLS的网络链路无防火墙/安全组限制
- Kafka权限:若Kafka配置了ACL,需确保自托管IR服务器的IP或认证账号拥有目标Topic的读取权限
- ADLS权限:使用服务主体认证时,需为其分配ADLS容器的「存储Blob数据贡献者」角色,避免写入失败
- 增量加载优化:可通过ADF的触发器(如定时触发器)配合Kafka的偏移量记录,实现周期性增量读取
- 监控与排查:通过ADF的「监控」面板查看管道运行日志,若出现连接失败,优先检查IR状态、网络连通性和认证配置
内容的提问来源于stack exchange,提问作者Ram Ranjan
相关产品推荐
相关产品推荐

