前言

在实时计算领域,数据是持续不断产生的,不能等待批量处理。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/File)
Sink 数据汇(输出)
Transformation 转换操作(Map/Filter/KeyBy)
Window 窗口(滚动/滑动/会话)
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: reduce
    map_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;
}
```

---

三、编译和运行

```bash
gcc -o flink flink.c -lpthread -lm
./flink
```

---

四、Flink vs 本实现

特性 本实现 Flink
流处理 ✅ ✅
窗口 ✅ 基础 ✅ 丰富
状态管理 ✅ ✅
检查点 ❌ ✅
时间语义 ✅ ✅
事件时间 ✅ ✅
水印 ❌ ✅
背压 ❌ ✅

---

五、总结

通过这篇文章,你学会了:

· Flink的核心架构(JobManager + TaskManager)
· 数据流与算子(Map/Filter/KeyBy)
· 窗口类型(滚动/滑动/会话)
· 状态管理(键控状态)
· 时间语义(Event/Processing Time)

Flink是分布式流处理的经典实现。掌握它,你就理解了实时计算系统的核心设计。

下一篇预告:《从零实现一个分布式任务调度:Apache Airflow的核心设计》

---

评论区分享一下你用Flink处理过什么实时场景~

更多推荐