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

使用wasbs协议向Azure Blob写入Parquet文件报错的解决方案咨询

使用wasbs协议创建Azure Blob Parquet Sink的问题与解决方案

问题背景

部署基于Flink的Parquet Sink代码时,使用wasbs协议连接Azure Blob存储,抛出异常,询问该协议是否可行及对应的修改方案。

报错信息

Caused by: java.lang.UnsupportedOperationException: Recoverable writers on AzureBlob are only supported for ABFS 
        at org.apache.flink.fs.azurefs.AzureBlobRecoverableWriter.checkSupportedFSSchemes(AzureBlobRecoverableWriter.java:44) ~[?:?]                                                           
        at org.apache.flink.fs.azure.common.hadoop.HadoopRecoverableWriter.<init>(HadoopRecoverableWriter.java:57) ~[?:?]                                                                    
        at org.apache.flink.fs.azurefs.AzureBlobRecoverableWriter.<init>(AzureBlobRecoverableWriter.java:37) ~[?:?]                                                                          
        at org.apache.flink.fs.azurefs.AzureBlobFileSystem.createRecoverableWriter(AzureBlobFileSystem.java:44) ~[?:?]                                                                        
        at org.apache.flink.core.fs.PluginFileSystemFactory$ClassLoaderFixingFileSystem.createRecoverableWriter(PluginFileSystemFactory.java:134) ~[flink-dist-1.18-SNAPSHOT.jar:1.18-SNAPSHOT]  

现有代码实现

object ParquetSink {

  def parquetFileSink[A <: Message: ClassTag](
      assigner: A => String,
      config: Config
  )(implicit lc: LoggingConfigs): FileSink[A] = {
    val bucketAssigner = new BucketAssigner[A, String] {
      override def getBucketId(element: A, context: BucketAssigner.Context): String = {
        val path = assigner(element)
        logger.info(LogMessage(-1, s"Writing file to ${config.getString(baseDirKey)}/$path", "NA"))
        path
      }

      override def getSerializer: SimpleVersionedSerializer[String] = SimpleVersionedStringSerializer.INSTANCE
    }

    def builder(outFile: OutputFile): ParquetWriter[A] =
      new ParquetProtoWriters.ParquetProtoWriterBuilder(
        outFile,
        implicitly[ClassTag[A]].runtimeClass.asInstanceOf[Class[A]]
      ).withCompressionCodec(config.getCompression(compressionKey)).build()

    val parquetBuilder: ParquetBuilder[A] = path => builder(path)
    FileSink
      .forBulkFormat(
        new Path(s"wasbs://${config.getString(baseDirKey)}@${config.getString(accountNameKey)}.blob.core.windows.net"),
        new ParquetWriterFactory[A](parquetBuilder)
      )
      .withBucketAssigner(bucketAssigner)
      .withOutputFileConfig(
        OutputFileConfig
          .builder()
          .withPartSuffix(".parquet")
          .build()
      )
      .build()
  }
}

解决方案

核心结论

wasbs协议不支持Flink的可恢复写入(Recoverable Writer)机制,这是Flink Azure Blob文件系统实现的明确限制:仅ABFS(Azure Blob File System)协议支持可恢复写入,wasbs作为旧版Blob存储协议,无法兼容该特性。

方案一:切换到ABFS协议(推荐生产环境)

ABFS是Flink官方推荐的Azure Blob存储访问方式,支持可恢复写入,能保证Exactly-Once语义,操作步骤如下:

  1. 确保Azure存储账户启用分层命名空间:ABFS要求存储账户必须开启此特性(创建账户时选择"BlobStorage"类型并启用分层命名空间,或对现有账户升级)。
  2. 修改代码中的路径协议与端点:
    • 将wasbs://替换为abfs://(非加密连接)或abfss://(加密HTTPS连接)
    • 将域名从blob.core.windows.net改为dfs.core.windows.net(ABFS使用DFS端点)
      修改后的路径构造代码:
    new Path(s"abfss://${config.getString(baseDirKey)}@${config.getString(accountNameKey)}.dfs.core.windows.net")
    
  3. 验证依赖配置:确保Flink集群已包含flink-azure-fs-hadoop依赖包(通常Flink发行版已内置,若缺失需手动添加)。

方案二:禁用可恢复写入(仅临时/非关键场景)

若无法切换到ABFS,可通过禁用FileSink的可恢复特性绕过限制,但会丢失Exactly-Once语义,故障恢复时可能出现数据重复或丢失,仅适合非生产环境测试:
在FileSink构建链中添加disableBucketCheckpoints()方法,禁用检查点驱动的桶提交机制:

FileSink
  .forBulkFormat(
    new Path(s"wasbs://${config.getString(baseDirKey)}@${config.getString(accountNameKey)}.blob.core.windows.net"),
    new ParquetWriterFactory[A](parquetBuilder)
  )
  .withBucketAssigner(bucketAssigner)
  .withOutputFileConfig(
    OutputFileConfig
      .builder()
      .withPartSuffix(".parquet")
      .build()
  )
  .disableBucketCheckpoints() // 禁用可恢复写入,放弃Exactly-Once语义
  .build()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 15:32:04