集成

Flink

师成师成· 更新于 2026-09-28· 阅读 8 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Apache Flink 是一个强大的开源分布式处理框架,专为在任意规模下对有界和无界数据流进行有状态计算而设计。它支持高吞吐、低延迟且容错的处理,同时提供弹性伸缩能力,可在数千个内核上每秒处理数百万个事件。

Apache Flink 可以使用 Apache Ozone 进行数据的读取与写入,也可以用它来存储关键的运维组件,例如应用状态的检查点(checkpoint)和保存点(savepoint)。

快速开始

本教程介绍如何使用 S3 Gateway 和 Docker Compose,让 Apache Flink 连接到 Apache Ozone。

快速开始环境

  • 非安全模式的 Ozone 和 Flink 集群。
  • Ozone S3G 启用了路径风格访问。若要启用虚拟主机风格寻址,请参见此处。
  • Flink 通过 S3 Gateway 访问 Ozone。

第 1 步 — 下载 Ozone 的 docker-compose.yaml

首先,获取 Ozone 的示例 Docker Compose 配置,并将其保存为 docker-compose.yaml:

curl -O https://raw.githubusercontent.com/apache/ozone-docker/refs/heads/latest/docker-compose.yaml

编辑 docker-compose.yaml:

将最后两个 SCM safemode 配置追加到 x-common-config: 部分,以便仅使用单个 Datanode 即可启动。

x-common-config:
   &common-config
   ...
   no_proxy: "om,recon,scm,s3g,localhost,127.0.0.1"
   OZONE-SITE.XML_hdds.scm.safemode.min.datanode: "1"
   OZONE-SITE.XML_hdds.scm.safemode.healthy.pipeline.pct: "0"

详情请参阅 Docker 快速开始页面。

第 2 步 — 为 Flink 创建 docker-compose-flink.yml

services:
  jobmanager:
    image: flink:scala_2.12-java17
    command: >
      bash -c "mkdir -p /opt/flink/plugins/s3-fs-hadoop &&
      cp /opt/flink/opt/flink-s3-fs-hadoop-*.jar /opt/flink/plugins/s3-fs-hadoop/ &&
      /docker-entrypoint.sh jobmanager"
    ports:
      - "8081:8081"
    environment:
      AWS_ACCESS_KEY_ID: ozone
      AWS_SECRET_ACCESS_KEY: ozone
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: jobmanager
        fs.s3a.endpoint: http://s3g:9878
        fs.s3a.path.style.access: true
        fs.s3a.connection.ssl.enabled: false
        fs.s3a.access.key: ozone
        fs.s3a.secret.key: ozone

  taskmanager:
    image: flink:scala_2.12-java17
    command: >
      bash -c "mkdir -p /opt/flink/plugins/s3-fs-hadoop &&
      cp /opt/flink/opt/flink-s3-fs-hadoop-*.jar /opt/flink/plugins/s3-fs-hadoop/ &&
      /docker-entrypoint.sh taskmanager"
    depends_on:
      - jobmanager
    environment:
      AWS_ACCESS_KEY_ID: ozone
      AWS_SECRET_ACCESS_KEY: ozone
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: jobmanager
        taskmanager.numberOfTaskSlots: 4
        fs.s3a.endpoint: http://s3g:9878
        fs.s3a.path.style.access: true
        fs.s3a.connection.ssl.enabled: false
        fs.s3a.access.key: ozone
        fs.s3a.secret.key: ozone

第 3 步 — 同时启动 Flink 和 Ozone

当用于 Ozone 的 docker-compose.yaml 和用于 Flink 的 docker-compose-flink.yml 位于同一目录中时,你可以使用以下命令同时启动这两个服务,并让它们共享同一网络:

export COMPOSE_FILE=docker-compose.yaml:docker-compose-flink.yml
docker compose up -d

验证容器正在运行:

docker ps

步骤 4 — 创建 Ozone 存储桶

你需要连接到 Ozone(例如 s3g)以创建 OBS 存储桶:

docker compose exec -it s3g ozone sh bucket create s3v/bucket1 -l obs

第 5 步 — 启动 Flink SQL 客户端

docker compose exec -it jobmanager ./bin/sql-client.sh

抱歉,没有收到需要翻译的 Markdown 文本内容。请提供完整的 Markdown 源文,我会按要求保留原有结构与语法,仅翻译其中的文字。

Flink SQL>

步骤 6 — 创建并查询由 Ozone S3 支持的表

重要:必须使用 BATCH 模式,否则分段上传会失败。

SET 'execution.runtime-mode' = 'BATCH';

CREATE TABLE ozone_sink (
  id STRING,
  ts TIMESTAMP(3)
) WITH (
  'connector' = 'filesystem',
  'path' = 's3a://bucket1/ozone_sink/',
  'format' = 'csv'
);

写入数据:

INSERT INTO ozone_sink VALUES ('hello', CURRENT_TIMESTAMP);

查询它:

SELECT * FROM ozone_sink;

如果一切正常,说明 Flink 已成功通过 S3 读写 Ozone。

第 7 步 — 在 Web UI 中检查 Flink 作业状态

打开浏览器:

http://localhost:8081/

在这里你可以:

  • 查看正在运行和已完成的作业
  • 检查 TaskManager
  • 以可视化方式调试故障

一旦出现问题,这里是首先应该查看的地方。

关键要点(重要)

  • Flink 的 Docker 镜像并未启用 S3
  • S3 插件必须同时存在于 JobManager(JM)和 TaskManager(TM)中
  • 应使用组合式的 Docker Compose 文件(COMPOSE_FILE)来启动 Flink 和 Ozone,以确保二者共享同一网络。
  • 使用 flink-s3-fs-hadoop 时,务必使用 s3a://
  • 访问 http://localhost:8081/ 以确认作业正在运行
  • Flink SQL 必须以批处理模式运行,以避免向 Ozone 的分段上传失败。请使用 SET 'execution.runtime-mode' = 'BATCH';

评论

登录后参与评论

正在加载评论…