Skip to content

Latest commit

 

History

History
85 lines (68 loc) · 2.65 KB

File metadata and controls

85 lines (68 loc) · 2.65 KB
name flink-dev
description Apache Flink 和 Flink CDC 开发专家。当用户提到 Flink 作业开发、CDC 数据同步、流处理、实时计算、DataStream API、状态管理、窗口操作、Flink 连接器开发、Pipeline 配置、Flink 测试时使用此技能。即使用户没有明确说"使用 flink-dev skill",只要涉及 Flink 相关开发都应该使用。

Flink 开发专家

协助 Apache Flink 和 Flink CDC 的开发、测试和部署。

核心能力

Flink 作业开发

  • DataStream API(Source、转换、Sink)
  • 状态管理(ValueState、ListState、MapState)
  • 窗口操作(滚动、滑动、会话窗口)
  • 水印和事件时间处理
  • 异步 I/O、广播状态、双流 Join

Flink CDC 开发

  • Pipeline YAML 配置(Transform、Route)
  • 事件模型(DataChangeEvent、SchemaChangeEvent)
  • 自定义连接器开发(DataSource、DataSink、Factory)
  • Schema 演化和全量增量同步

测试和部署

  • 单元测试和集成测试(Testcontainers)
  • 性能测试(吞吐量、延迟)
  • 打包部署和 Savepoint 管理

使用方式

当需要详细参考时,查看:

  • references/flink-dev-guide.md - Flink 开发完整指南
  • references/flink-cdc-guide.md - Flink CDC 开发指南
  • references/flink-connector-dev.md - 连接器开发完整指南
  • references/flink-plugins.md - 插件机制和扩展点
  • references/flink-builtin-tools.md - 内置工具和开发规范
  • references/flink-ha.md - 高可用配置(ZooKeeper/Kubernetes/Standalone)
  • references/flink-test-guide.md - 测试规范和最佳实践
  • references/flink-test-categories.md - 测试分类
  • references/flink-architecture.md - 源码架构要点
  • references/flink-best-practices.md - 实战要点和常见陷阱

快速示例

Flink DataStream 作业

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60000);

DataStream<String> stream = env.fromSource(kafkaSource, watermarkStrategy, "Kafka");
stream
    .map(String::toUpperCase)
    .keyBy(s -> s)
    .window(TumblingEventTimeWindows.of(Duration.ofMinutes(5)))
    .reduce((a, b) -> a + b)
    .sinkTo(sink);

env.execute("My Job");

Flink CDC Pipeline

source:
  type: mysql
  hostname: localhost
  port: 3306
  tables: db.\.*

sink:
  type: doris
  fenodes: 127.0.0.1:8030

pipeline:
  parallelism: 4
  schema.change.behavior: evolve

开发规范

  • 使用 Google Java Format (AOSP)
  • Commit 格式: [FLINK-XXXX][module] 描述
  • 测试命名: *Test.java (单元), *ITCase.java (集成)
  • 状态后端: RocksDB 用于大状态
  • Checkpoint 间隔: 60s 起步