- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】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 | 端点域名 | 适用存储类型 | 加密传输 |
|---|---|---|---|---|
| WASB | wasb:// | *.blob.core.windows.net | Azure Blob Storage | 否 |
| WASB(SSL) | wasbs:// | *.blob.core.windows.net | Azure Blob Storage | 是 |
| ABFS | abfs:// | *.dfs.core.windows.net | ADLS Gen2 存储账户 | 否 |
| ABFS(SSL) | abfss:// | *.dfs.core.windows.net | ADLS 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 文件系统 |
|---|---|---|
| AzureBlobStorageFSFactory | wasb | NativeAzureFileSystem |
| SecureAzureBlobStorageFSFactory | wasbs | NativeAzureFileSystem |
| AzureDataLakeStoreGen2FSFactory | abfs | hadoop-azurebfs的AzureBlobFileSystem |
| SecureAzureDataLakeStoreGen2FSFactory | abfss | hadoop-azurebfs的AzureBlobFileSystem |
3.3 Shaded JAR 的依赖与类重定位
flink-azure-fs-hadoop是一个shaded(重定位)插件包:pom.xml 依赖hadoop-azure(并排除其传递的hadoop-common、reload4j等,避免与 Flink 自身依赖冲突),同时通过maven-shade-plugin将org.apache.flink.runtime.fs.hdfs与org.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 文件系统的支持遵循统一的插件工厂模式:
- Flink 运行时按 URI Scheme 查找
FileSystemFactory(SPI 加载,来源为插件 JAR 中的META-INF/services文件); - 命中工厂后调用
AbstractAzureFSFactory.configure(Configuration),将 Flink 配置经HadoopConfigLoader映射为 Hadoop 配置; create(URI)中实例化对应的 Hadoop 文件系统(NativeAzureFileSystem或AzureBlobFileSystem),用映射后的 Hadoop 配置执行initialize;- 将 Hadoop 文件系统包装进 Flink 的 AzureBlobFileSystem(继承自
HadoopFileSystem),并覆写createRecoverableWriter()返回AzureBlobRecoverableWriter,从而支持 Flink 的可恢复写入(RecoverableWriter)机制——这正是 StreamingFileSink、FileSink 以及 checkpoint 元数据写入能够对 Azure 对象存储做故障恢复的基础。
5.2 测试与验证
仓库中为该模块提供了单元测试与集成测试,可用于理解其行为约定:
- AzureBlobStorageFSFactoryTest 与 AzureDataLakeStoreGen2FSFactoryTest 验证工厂 Scheme 注册与配置加载;
- AzureBlobRecoverableWriterTest 与 AzureBlobFsRecoverableDataOutputStreamTest 覆盖可恢复写入与输出流的提交/恢复语义;
- AzureFileSystemBehaviorITCase 是针对真实 Azure 环境的文件系统行为集成测试。
六、常见问题与建议
- Scheme 选择:普通 Blob 账户用
wasb:///wasbs://;ADLS Gen2 账户优先用abfs:///abfss://(Azure 官方推荐),但注意abfs://不能访问非 Gen2 账户。 - 生产环境传输安全:涉及敏感数据时使用带 SSL 的
wasbs://或abfss://。 - 密钥安全:优先使用托管身份或环境变量 KeyProvider,避免将存储密钥明文写死在
flink-conf.yaml中;使用环境变量方式时确保所有节点进程均注入AZURE_STORAGE_KEY。 - 插件目录:JAR 必须放入独立的
plugins/<name>/子目录(每插件一个文件夹),不要散落在lib/下,以利用插件类加载器的隔离能力。 - checkpoint 目录:将
execution.checkpointing.dir或CheckpointingOptions.CHECKPOINTS_DIRECTORY指向 Azure 路径即可把对象存储用作 checkpoint 存储,需保证所有 TaskManager 具备相同的凭据配置。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
Pixelle-Video终极指南:如何用AI全自动制作专业短视频
Pixelle Video终极指南:如何用AI全自动制作专业短视频 在当今内容为王的时代,视频创作已经成为信息传播的核心方式。然而,传统的视频制作需要文案撰写、
人工智能AI 应用音视频媒体生成为什么选择gh_mirrors/ipd/IP_database?5大优势解析
为什么选择gh_mirrors/ipd/IP_database?5大优势解析 gh_mirrors/ipd/IP_database是一个全面的IP地址库项目,提
FiftyOne Enterprise 安装部署与云存储凭据配置实战指南
FiftyOne Enterprise 安装部署与云存储凭据配置实战指南 FiftyOne Enterprise 是 Voxel51 面向团队协作推出的企业版部
人工智能计算机视觉数据集数据可视化数据标注模型评测
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考