使用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语义,操作步骤如下:
- 确保Azure存储账户启用分层命名空间:ABFS要求存储账户必须开启此特性(创建账户时选择"BlobStorage"类型并启用分层命名空间,或对现有账户升级)。
- 修改代码中的路径协议与端点:
- 将
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") - 将
- 验证依赖配置:确保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
相关产品推荐
相关产品推荐

