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

通过Kafka API从Azure Databricks向Event Hub传输CSV流数据问题排查

排查思路与修复方案

流启动后输入速率为0,按以下优先级逐个排查修复:

  • 文件流源默认行为问题(最高发原因)
    spark.readStream.csv默认仅处理流启动后新上传到目标路径的文件,启动前已存在的CSV文件不会被扫描读取。
    修复方式:读取文件时追加配置,强制按文件修改时间从早到晚处理,包含路径下已有文件:
    streaming_df = spark.readStream \
      .option("header", "true") \
      .option("latestFirst", "false") \
      .option("maxFilesPerTrigger", 10) # 按需求调整每批处理文件数,控制流速率
      .schema(location_schema) \
      .csv(f"{path}")
    
  • Event Hubs Kafka兼容层参数配置错误
    现有配置存在三个典型错漏:
    1. kafka.bootstrap.servers格式错误:正确值应为<你的Event Hubs命名空间名>.servicebus.windows.net:9093,不能直接填入连接字符串的Endpoint片段
    2. kafka.sasl.jaas.config格式错误:不能直接填入Event Hubs连接字符串,正确格式如下,注意username必须固定为$ConnectionString,password为你拿到的完整Event Hubs连接字符串,配置末尾必须加英文分号:
      org.apache.kafka.common.security.plain.PlainLoginModule required username="$ConnectionString" password="Endpoint=sb://XXXX.servicebus.windows.net/;SharedAccessKeyName=XXXX;SharedAccessKey=XXX=;EntityPath=XXXX";
      
    3. 缺少超时配置:追加两个超时参数避免连接被静默断开无报错:
      .option("kafka.request.timeout.ms", "60000") \
      .option("kafka.session.timeout.ms", "60000")
      
  • 检查点路径非法
    当前使用的./checkpoint是集群本地相对路径,Databricks运行环境下本地路径无持久化权限,会导致检查点写入失败、流初始化异常但无显性报错。
    修复方式:替换为DBFS绝对路径,例如dbfs:/checkpoints/event_hub_csv_stream/,每次修改流逻辑后清空对应检查点目录再重启,避免旧状态干扰。
  • 写入数据不符合Kafka格式要求
    Kafka Sink要求待写入DataFrame必须存在名为value的列(可选key列),其余列会被直接忽略。当前直接select("*")写入,如果CSV数据中无value列,写入操作会静默失败。
    修复方式:写入前将待发送数据转换为value列,例如将整行数据转为JSON字符串:
    from pyspark.sql.functions import to_json, struct
    stream_write_df = streaming_df.select(to_json(struct("*")).alias("value"))
    # 后续将stream_write_df传入写入方法即可
    

修复后启动流,可直接在Driver日志中搜索KafkaProducer关键词,排查是否存在认证失败、连接拒绝类报错;同时查看流指标的numInputRows字段,该值大于0即代表文件已被正常读取。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 13:15:24