Flink 集成 Azure Blob 存储:wasb/abfs 访问、插件部署与凭据配置实战指南
2026/9/20 14:38:29 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

Azure Blob 存储是 Microsoft 托管的对象存储服务,在 Flink 中可像普通文件一样被用于读取和写入数据,也可以作为流式作业的 State Backend 的存储载体(例如将 checkpoint 目录指向 Azure 路径)。本文基于 Apache Flink 仓库中 docs/content.zh/docs/deployment/filesystems/azure.md 展开,结合 flink-azure-fs-hadoop 模块源码,系统讲解 Flink 如何通过wasb:///wasbs://abfs:///abfss://四种 URI Scheme 访问 Azure Blob 存储,包括路径格式、插件化部署方式、凭据配置(WASB 存储密钥、环境变量 KeyProvider、ABFS 存储密钥与托管身份)以及底层实现原理,帮助你在一线生产环境中正确、安全地接入 Azure 存储。

一、访问协议与路径格式

Flink 支持两种 Azure Blob 存储访问协议:WASB(Windows Azure Storage Blob)与ABFS(Azure Blob File System)。前者由 Hadoop 的hadoop-azure模块提供,后者由hadoop-azure-datalake模块提供,二者在 Flink 中分别对应不同的文件系统工厂类。

协议Scheme端点域名适用存储类型加密传输
WASBwasb://*.blob.core.windows.netAzure Blob Storage
WASB(SSL)wasbs://*.blob.core.windows.netAzure Blob Storage
ABFSabfs://*.dfs.core.windows.netADLS Gen2 存储账户
ABFS(SSL)abfss://*.dfs.core.windows.netADLS Gen2 存储账户

提示:Azure 官方推荐使用abfs://访问 ADLS Gen2 存储账户,尽管wasb://通过向后兼容也能工作。

警告abfs://只能用于访问 ADLS Gen2 存储账户,普通 Blob 存储账户请使用wasb://。如何识别 ADLS Gen2 存储账户请参阅 Azure 官方文档。

在 Flink 中,Azure Blob 存储对象可以像普通文件一样被引用,路径格式如下:

// WASB 非加密访问 wasb://<your-container>@$<your-azure-account>.blob.core.windows.net/<object-path> // WASB SSL 加密访问 wasbs://<your-container>@$<your-azure-account>.blob.core.windows.net/<object-path> // ABFS 非加密访问 abfs://<your-container>@$<your-azure-account>.dfs.core.windows.net/<object-path> // ABFS SSL 加密访问 abfss://<your-container>@$<your-azure-account>.dfs.core.windows.net/<object-path>

路径由三部分组成:<your-container>(容器名)、$<your-azure-account>(Azure 存储账户名)、<object-path>(对象在容器内的路径)。注意账户名前保留了$符号,这是 WASB/ABFS URI 的标准语法。

二、在 Flink 作业中使用 Azure Blob 存储

得到上述路径后,可以把它当作本地文件路径直接传入 DataStream API。以下代码演示读取、写入以及将 Azure 路径用作 checkpoint 存储:

// 读取 Azure Blob 存储 env.readTextFile("wasb://<your-container>@$<your-azure-account>.blob.core.windows.net/<object-path>"); // 写入 Azure Blob 存储 stream.writeAsText("wasb://<your-container>@$<your-azure-account>.blob.core.windows.net/<object-path>"); // 将 Azure Blob 存储用作 checkpoint storage Configuration config = new Configuration(); config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem"); config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "wasb://<your-container>@$<your-azure-account>.blob.core.windows.net/<object-path>"); env.configure(config);

其中 checkpoint 目录也可以在 Flink 配置文件 中通过execution.checkpointing.dir全局指定,效果与上述CheckpointingOptions.CHECKPOINTS_DIRECTORY一致——该选项指定了所有 State Backend 写 Checkpoint 数据和元数据文件的目录(参见 State Backends)。对于流式作业,通常建议checkpoint storage使用filesystem,并将目录指向对象存储,以便在 TaskManager 故障或集群重启后从 Azure 上恢复状态。

三、部署 flink-azure-fs-hadoop 插件

3.1 插件化加载机制

从 Flink 1.9 开始,文件系统采用插件机制加载(参见 文件系统插件 与 文件系统总览):文件系统 Factory 类由专用 Java 类加载器加载,避免与其他类或 Flink 组件冲突。因此使用 Azure 文件系统前,需要把对应的 JAR 从opt目录复制到发行版的plugins目录:

mkdir ./plugins/azure-fs-hadoop cp ./opt/flink-azure-fs-hadoop-{{< version >}}.jar ./plugins/azure-fs-hadoop/

{{< version >}}替换为你所用 Flink 发行版的版本号(如1.18.0)。该 JAR 由构建流程自动产出:在 flink-dist/src/main/assemblies/opt.xml 中,../flink-filesystems/flink-azure-fs-hadoop/target/flink-azure-fs-hadoop-${project.version}.jar被打包进发行版opt目录。

3.2 插件注册了哪些 Scheme

flink-azure-fs-hadoop通过 ServiceLoader 机制注册了四个文件系统工厂,见 META-INF/services/org.apache.flink.core.fs.FileSystemFactory:

org.apache.flink.fs.azurefs.AzureBlobStorageFSFactory org.apache.flink.fs.azurefs.SecureAzureBlobStorageFSFactory org.apache.flink.fs.azurefs.AzureDataLakeStoreGen2FSFactory org.apache.flink.fs.azurefs.SecureAzureDataLakeStoreGen2FSFactory

四个工厂类分别声明了对应 Scheme,并基于 Hadoop 的NativeAzureFileSystem(WASB)或org.apache.hadoop.fs.azurebfs.AzureBlobFileSystem(ABFS)创建底层文件系统:

工厂类Scheme底层 Hadoop 文件系统
AzureBlobStorageFSFactorywasbNativeAzureFileSystem
SecureAzureBlobStorageFSFactorywasbsNativeAzureFileSystem
AzureDataLakeStoreGen2FSFactoryabfshadoop-azurebfsAzureBlobFileSystem
SecureAzureDataLakeStoreGen2FSFactoryabfsshadoop-azurebfsAzureBlobFileSystem

3.3 Shaded JAR 的依赖与类重定位

flink-azure-fs-hadoop是一个shaded(重定位)插件包:pom.xml 依赖hadoop-azure(并排除其传递的hadoop-commonreload4j等,避免与 Flink 自身依赖冲突),同时通过maven-shade-pluginorg.apache.flink.runtime.fs.hdfsorg.apache.flink.runtime.util重定位到org.apache.flink.fs.azure.common.*命名空间,强制这些适配类使用插件专属类加载器,从而避免与其他文件系统插件的类冲突。这也是官方文档要求"将插件 JAR 放入独立插件目录"而非直接丢进lib的原因。

四、凭据配置

4.1 WASB:Hadoop 配置方式

Hadoop 的 WASB 文件系统支持通过 Hadoop 配置来配置凭据。为方便起见,Flink 会把所有以fs.azure.前缀开头的 Flink 配置转发到文件系统的 Hadoop 配置中。这一点可在源码中得到印证:AbstractAzureFSFactory 中定义了FLINK_CONFIG_PREFIXES = {"fs.azure.", "azure."}HADOOP_CONFIG_PREFIX = "fs.azure.",通过HadoopConfigLoader完成配置映射;在create(URI)时加载 Hadoop 配置并调用fs.initialize(fsUri, hadoopConfig)初始化底层文件系统。

因此,可以在 Flink 配置文件(conf/flink-conf.yaml)中直接配置 Azure Blob 存储密钥:

fs.azure.account.key.<account_name>.blob.core.windows.net: <azure_storage_key>

其中<account_name>是存储账户名,<azure_storage_key>是账户的访问密钥(可在 Azure 门户的存储账户 → 访问密钥中获取)。

4.2 WASB:从环境变量读取密钥

为避免把密钥明文写入配置文件,可以将文件系统配置为从环境变量AZURE_STORAGE_KEY读取密钥。在 Flink 配置文件中设置:

fs.azure.account.keyprovider.<account_name>.blob.core.windows.net: org.apache.flink.fs.azurefs.EnvironmentVariableKeyProvider

对应的 EnvironmentVariableKeyProvider 实现了 Hadoop 的KeyProvider接口:其getStorageAccountKey方法从环境变量AZURE_STORAGE_KEY(常量AZURE_STORAGE_KEY_ENV_VARIABLE)读取密钥并返回;若环境变量未设置,则抛出KeyProviderException,提示"AZURE_STORAGE_KEY" not set。使用该方式时,需确保运行 Flink 的进程(JobManager 与 TaskManager)都能访问到此环境变量。

4.3 ABFS:存储密钥方式(不鼓励)

Hadoop 的 ABFS 文件系统支持多种认证方式。一种简单但不鼓励的做法是直接使用存储密钥:在 Flink 配置文件中配置

fs.azure.account.key.<account_name>.dfs.core.windows.net: <azure_storage_key>

注意此处端点为dfs.core.windows.net,与 WASB 的blob.core.windows.net不同。

4.4 ABFS:Azure 托管身份(推荐)

Azure 官方推荐使用Azure 托管身份(Managed Identity)通过abfs://访问 ADLS Gen2 存储账户,从而避免在配置中暴露任何密钥。其配置细节属于 Azure 平台侧操作(如为集群开启系统/用户分配托管身份,并在存储账户的 IAM 中授予相应 RBAC 角色),具体步骤请参阅 Azure 文档中关于托管身份的支持服务列表——部署在支持托管身份的 Azure 服务(如 Azure VM、AKS、HDInsight 等)上的 Flink 集群都可以利用该能力,Flink 侧无需再配置存储密钥。

五、底层实现原理与验证

5.1 工厂模式的加载链路

从源码结构看,Flink 对 Azure 文件系统的支持遵循统一的插件工厂模式:

  1. Flink 运行时按 URI Scheme 查找FileSystemFactory(SPI 加载,来源为插件 JAR 中的META-INF/services文件);
  2. 命中工厂后调用AbstractAzureFSFactory.configure(Configuration),将 Flink 配置经HadoopConfigLoader映射为 Hadoop 配置;
  3. create(URI)中实例化对应的 Hadoop 文件系统(NativeAzureFileSystemAzureBlobFileSystem),用映射后的 Hadoop 配置执行initialize
  4. 将 Hadoop 文件系统包装进 Flink 的 AzureBlobFileSystem(继承自HadoopFileSystem),并覆写createRecoverableWriter()返回AzureBlobRecoverableWriter,从而支持 Flink 的可恢复写入(RecoverableWriter)机制——这正是 StreamingFileSink、FileSink 以及 checkpoint 元数据写入能够对 Azure 对象存储做故障恢复的基础。

5.2 测试与验证

仓库中为该模块提供了单元测试与集成测试,可用于理解其行为约定:

  • AzureBlobStorageFSFactoryTest 与 AzureDataLakeStoreGen2FSFactoryTest 验证工厂 Scheme 注册与配置加载;
  • AzureBlobRecoverableWriterTest 与 AzureBlobFsRecoverableDataOutputStreamTest 覆盖可恢复写入与输出流的提交/恢复语义;
  • AzureFileSystemBehaviorITCase 是针对真实 Azure 环境的文件系统行为集成测试。

六、常见问题与建议

  1. Scheme 选择:普通 Blob 账户用wasb:///wasbs://;ADLS Gen2 账户优先用abfs:///abfss://(Azure 官方推荐),但注意abfs://不能访问非 Gen2 账户。
  2. 生产环境传输安全:涉及敏感数据时使用带 SSL 的wasbs://abfss://
  3. 密钥安全:优先使用托管身份或环境变量 KeyProvider,避免将存储密钥明文写死在flink-conf.yaml中;使用环境变量方式时确保所有节点进程均注入AZURE_STORAGE_KEY
  4. 插件目录:JAR 必须放入独立的plugins/<name>/子目录(每插件一个文件夹),不要散落在lib/下,以利用插件类加载器的隔离能力。
  5. checkpoint 目录:将execution.checkpointing.dirCheckpointingOptions.CHECKPOINTS_DIRECTORY指向 Azure 路径即可把对象存储用作 checkpoint 存储,需保证所有 TaskManager 具备相同的凭据配置。
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

上一篇:Cube Sandbox 持久化存储实战:Host Mount 主机挂载完整指南
下一篇:在Android手机上玩Minecraft Java版的终极指南:MCinaBox启动器完全解析

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询