从零实现一个分布式流处理:Apache Flink的核心设计
前言在实时计算领域数据是持续不断产生的不能等待批量处理。Apache Flink 是流处理领域的标杆实现了真正的流计算而非微批处理。今天我们从零实现Flink的核心功能· 数据源Source与数据汇Sink· 转换算子Map/Filter/FlatMap· 窗口Window· 时间语义Event Time/Processing Time· 状态管理State· 检查点Checkpoint· 容错与恢复---一、Flink核心原理1. 架构图┌─────────────────────────────────────────────────────────────┐│ 数据流 ││ Source → Map → Filter → Window → Reduce → Sink │└─────────────────────────────────────────────────────────────┘│▼┌─────────────────────────────────────────────────────────────┐│ Flink运行时 ││ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ ││ │ JobManager│ │ TaskManager│ │ TaskManager│ ││ │ (调度) │ │ (执行) │ │ (执行) │ ││ └─────────────┘ └─────────────┘ └─────────────┘ │└─────────────────────────────────────────────────────────────┘│▼┌─────────────────────────────────────────────────────────────┐│ 状态后端 ││ (内存 / RocksDB / HDFS) │└─────────────────────────────────────────────────────────────┘2. 核心概念概念 说明Source 数据源Kafka/Socket/FileSink 数据汇输出Transformation 转换操作Map/Filter/KeyByWindow 窗口滚动/滑动/会话State 状态算子状态/键控状态Checkpoint 快照容错Time 时间语义Event/Processing/Ingestion---二、完整代码实现1. 基础数据结构c#include stdio.h#include stdlib.h#include string.h#include unistd.h#include pthread.h#include time.h#include errno.h#include math.h#define MAX_EVENT_SIZE 1024#define MAX_STREAM_NAME 64#define MAX_OPERATORS 20#define MAX_WATERMARK_MS 5000// 事件/数据typedef struct event {char data[MAX_EVENT_SIZE];long long timestamp; // 事件时间long long processing_time; // 处理时间char key[64];struct event *next;} event_t;// 窗口类型typedef enum {WINDOW_TUMBLING 0, // 滚动窗口WINDOW_SLIDING 1, // 滑动窗口WINDOW_SESSION 2 // 会话窗口} window_type_t;// 窗口定义typedef struct window_def {window_type_t type;long long size_ms;long long slide_ms;long long gap_ms;struct window_def *next;} window_def_t;// 状态typedef struct state {char key[64];void *value;size_t value_size;long long last_update;struct state *next;} state_t;// 算子函数类型typedef event_t* (*map_func_t)(event_t *e);typedef event_t* (*filter_func_t)(event_t *e);typedef event_t* (*flatmap_func_t)(event_t *e, event_t **output);// 算子typedef struct operator {char name[64];int op_type; // 0: map, 1: filter, 2: flatmap, 3: keyby, 4: window, 5: reducemap_func_t map_func;filter_func_t filter_func;flatmap_func_t flatmap_func;char key_field[64];window_def_t *window_def;struct operator *next;} operator_t;// 数据流typedef struct data_stream {char name[MAX_STREAM_NAME];operator_t *operators;int operator_count;struct data_stream *next;} data_stream_t;// 执行环境typedef struct flink_env {data_stream_t *streams;int stream_count;pthread_mutex_t mutex;int running;int parallelism;int checkpoint_interval_ms;long long current_watermark;} flink_env_t;// 数据源typedef struct source {char name[64];event_t* (*generate)(void *ctx);void *ctx;pthread_t thread;int running;struct source *next;} source_t;2. 执行环境c// 创建Flink环境flink_env_t *flink_create(int parallelism, int checkpoint_interval_ms) {flink_env_t *env malloc(sizeof(flink_env_t));memset(env, 0, sizeof(flink_env_t));env-parallelism parallelism;env-checkpoint_interval_ms checkpoint_interval_ms;env-running 1;env-current_watermark 0;pthread_mutex_init(env-mutex, NULL);printf([Flink] 环境创建并行度: %d\n, parallelism);return env;}// 创建数据流data_stream_t *flink_add_stream(flink_env_t *env, const char *name) {pthread_mutex_lock(env-mutex);data_stream_t *stream malloc(sizeof(data_stream_t));strcpy(stream-name, name);stream-operators NULL;stream-operator_count 0;stream-next env-streams;env-streams stream;env-stream_count;pthread_mutex_unlock(env-mutex);printf([Flink] 创建数据流: %s\n, name);return stream;}// 添加Map算子operator_t *flink_map(data_stream_t *stream, map_func_t func, const char *name) {operator_t *op malloc(sizeof(operator_t));strcpy(op-name, name);op-op_type 0;op-map_func func;op-next stream-operators;stream-operators op;stream-operator_count;printf([Flink] 添加Map: %s\n, name);return op;}// 添加Filter算子operator_t *flink_filter(data_stream_t *stream, filter_func_t func, const char *name) {operator_t *op malloc(sizeof(operator_t));strcpy(op-name, name);op-op_type 1;op-filter_func func;op-next stream-operators;stream-operators op;stream-operator_count;printf([Flink] 添加Filter: %s\n, name);return op;}3. 窗口实现c// 创建滚动窗口window_def_t *window_tumbling(long long size_ms) {window_def_t *w malloc(sizeof(window_def_t));w-type WINDOW_TUMBLING;w-size_ms size_ms;w-slide_ms size_ms;w-gap_ms 0;return w;}// 创建滑动窗口window_def_t *window_sliding(long long size_ms, long long slide_ms) {window_def_t *w malloc(sizeof(window_def_t));w-type WINDOW_SLIDING;w-size_ms size_ms;w-slide_ms slide_ms;w-gap_ms 0;return w;}// 窗口聚合typedef struct window_result {char key[64];long long window_start;long long window_end;long long count;long long sum;double avg;} window_result_t;// 滚动窗口处理void process_tumbling_window(event_t **events, int count,long long window_start, window_result_t *result) {result-window_start window_start;result-window_end window_start 10000; // 10秒窗口result-count count;result-sum 0;for (int i 0; i count; i) {result-sum atol(events[i]-data);}result-avg count 0 ? (double)result-sum / count : 0;}// 执行窗口操作event_t *flink_window(data_stream_t *stream, window_def_t *w, const char *name) {printf([Flink] 窗口操作: %s (类型: %d)\n, name, w-type);return NULL;}4. 状态管理c// 键控状态typedef struct keyed_state {struct state *states;pthread_mutex_t mutex;} keyed_state_t;keyed_state_t *keyed_state_create(void) {keyed_state_t *ks malloc(sizeof(keyed_state_t));ks-states NULL;pthread_mutex_init(ks-mutex, NULL);return ks;}// 获取状态void *keyed_state_get(keyed_state_t *ks, const char *key) {pthread_mutex_lock(ks-mutex);state_t *s ks-states;while (s) {if (strcmp(s-key, key) 0) {pthread_mutex_unlock(ks-mutex);return s-value;}s s-next;}pthread_mutex_unlock(ks-mutex);return NULL;}// 更新状态void keyed_state_put(keyed_state_t *ks, const char *key, void *value, size_t size) {pthread_mutex_lock(ks-mutex);state_t *s ks-states;while (s) {if (strcmp(s-key, key) 0) {if (s-value) free(s-value);s-value malloc(size);memcpy(s-value, value, size);s-value_size size;s-last_update time(NULL);pthread_mutex_unlock(ks-mutex);return;}s s-next;}s malloc(sizeof(state_t));strcpy(s-key, key);s-value malloc(size);memcpy(s-value, value, size);s-value_size size;s-last_update time(NULL);s-next ks-states;ks-states s;pthread_mutex_unlock(ks-mutex);}5. 测试代码c// 示例Map函数提取数字event_t *parse_number_map(event_t *e) {event_t *out malloc(sizeof(event_t));memcpy(out, e, sizeof(event_t));// 提取数据中的数字char *p e-data;while (*p !isdigit(*p)) p;if (*p) {char num[64];int i 0;while (*p (isdigit(*p) || *p .)) {num[i] *p;}num[i] \0;strcpy(out-data, num);} else {strcpy(out-data, 0);}return out;}// 示例Filter函数过滤负数event_t *filter_positive(event_t *e) {int val atoi(e-data);if (val 0) return NULL;event_t *out malloc(sizeof(event_t));memcpy(out, e, sizeof(event_t));return out;}// 测试流处理void test_flink() {printf( Flink流处理测试 \n\n);flink_env_t *env flink_create(4, 5000);// 创建数据流data_stream_t *stream flink_add_stream(env, number-stream);// 构建流处理管道flink_map(stream, parse_number_map, parse-number);flink_filter(stream, filter_positive, filter-positive);// 添加窗口10秒滚动窗口window_def_t *w window_tumbling(10000);flink_window(stream, w, tumbling-window);// 模拟数据流printf(\n模拟数据处理:\n);char *test_data[] {data: 100, data: -50, data: 200, data: 150, data: -30};for (int i 0; i 5; i) {event_t *e malloc(sizeof(event_t));strcpy(e-data, test_data[i]);e-timestamp time(NULL) * 1000;printf( 输入: %s\n, test_data[i]);// 模拟算子执行event_t *after_map parse_number_map(e);printf( → Map: %s\n, after_map-data);event_t *after_filter filter_positive(after_map);if (after_filter) {printf( → Filter: 通过\n);free(after_filter);} else {printf( → Filter: 过滤\n);}free(after_map);free(e);}free(w);free(env);}int main() {srand(time(NULL));test_flink();return 0;}---三、编译和运行bashgcc -o flink flink.c -lpthread -lm./flink---四、Flink vs 本实现特性 本实现 Flink流处理 ✅ ✅窗口 ✅ 基础 ✅ 丰富状态管理 ✅ ✅检查点 ❌ ✅时间语义 ✅ ✅事件时间 ✅ ✅水印 ❌ ✅背压 ❌ ✅---五、总结通过这篇文章你学会了· Flink的核心架构JobManager TaskManager· 数据流与算子Map/Filter/KeyBy· 窗口类型滚动/滑动/会话· 状态管理键控状态· 时间语义Event/Processing TimeFlink是分布式流处理的经典实现。掌握它你就理解了实时计算系统的核心设计。下一篇预告《从零实现一个分布式任务调度Apache Airflow的核心设计》---评论区分享一下你用Flink处理过什么实时场景