用 Kafka Connect JDBC Sink 给 MySQL 写吞吐做基线测试con/connect 仓库 jdbc-sink 基准实践【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect导读本文以 internal/impl/mysql/bench/mysql-write/jdbc-sink/README.md 为核心完整拆解 con/connect 仓库中“Confluent JDBC Sink 连接器从 Kafka 写 MySQL”的吞吐基准benchmark方案。该基准用于与 Redpanda Connect 的sql_insert输出做对比基线衡量“流处理引擎写数据库”这条链路的真实上限。读完本文你将掌握JDBC Sink 基准的完整运行流程起环境、造数、跑测、清理、连接器关键配置tasks.max、batch.size、consumer.override 系列参数的调优含义以及如何在仓库内复现同一套对比实验并读懂 docs/benchmark-results/mysql-cdc.md 中的结果表。一、基准是什么为何要给“写 MySQL”单独测一条链路在流处理场景里向关系型数据库写入通常是整条管线的瓶颈与终点。Kafka Connect JDBC Sink 是社区最常用的“Kafka 主题 → 关系型数据库”落库方案之一因此本仓库把它当作对比基线用来回答两个问题一个成熟的、久经考验的 Kafka Connect 生态连接器把数据从 Kafka 写进 MySQL 能达到多高的稳定吞吐换成 Redpanda Connect 的sql_insert输出见 internal/impl/mysql/bench/mysql-write/rpcn/ 对应实现 output_sql_insert.go在更少的资源下是否更快整个 mysql-write 基准目录刻意把两条候选路径平铺开internal/impl/mysql/bench/mysql-write/ ├── jdbc-sink/ # Kafka Connect Confluent JDBC Sink本文主体 └── rpcn/ # Redpanda Connect kafka_franz → sql_insert两者共用同一套数据模型bench_events表、16 分区的bench-events主题保证对比公平。二、前置条件与整体架构2.1 前置条件原文档明确要求Docker 必须处于运行状态。所有基础设施Kafka、MySQL、Kafka Connect都通过 docker compose 拉起本机无需预装 Kafka 或 MySQL 客户端工具命令都走docker exec。2.2 服务组成docker-compose.yaml 定义了 5 个服务服务镜像作用kafkaconfluentinc/cp-kafka:7.7.8KRaft 模式单节点 Kafka无 ZooKeepercpus: 3mysqlmysql:8.0落库目标预建benchdb库与bench_events表kafka-connect基于 Dockerfile 构建装有 Confluent JDBC Sink 插件REST 端口 8083kafka-setupconfluentinc/cp-kafka:7.7.8一次性创建bench-events主题16 分区、复制因子 1mysql-setupmysql:8.0一次性创建目标表bench_events表结构由mysql-setup服务创建CREATE TABLE IF NOT EXISTS bench_events ( id BIGINT, category VARCHAR(50), value DOUBLE, ts BIGINT );2.3 关键部署细节Kafka Connect 镜像由 Dockerfile 现场构建在confluentinc/cp-kafka-connect-base:7.7.8基础上通过confluent-hub install安装confluentinc/kafka-connect-jdbc:10.7.14再手动下载 MySQL JDBC 驱动mysql-connector-j-8.0.33.jar该驱动并不随 hub 插件包分发必须单独安装否则连接器无法连接 MySQL。Kafka Connect worker 的 key/value 转换器在 compose 里声明为StringConverterJsonConverter且CONNECT_VALUE_CONVERTER_SCHEMAS_ENABLE: false。两个setup服务都依赖目标服务service_healthyKafka 用kafka-broker-api-versions探活MySQL 用mysqladmin ping探活Kafka Connect 用curl http://localhost:8083/探活。三、一条命令起环境task up环境编排全部收敛在 Taskfile.yaml仓库约定以tasktaskfile.dev作为命令入口。启动命令task up其实际行为对应up任务docker compose up -d --build until curl -sf http://localhost:8083/ /dev/null 21; do sleep 3; done echo Kafka Connect ready at http://localhost:8083--build会触发上文 Dockerfile 的镜像构建因此首次启动较慢直到 Kafka Connect 的 REST 端点可访问才算就绪。任务默认值均可通过环境变量覆盖KAFKA_BOOTSTRAP默认localhost:9092、CONNECT_URL默认http://localhost:8083、CONNECTOR_NAME默认mysql-jdbc-sink-bench、CONSUMER_GROUP默认connect-mysql-jdbc-sink-bench。四、造数向 bench-events 生产 1000 万条事件4.1 命令task bench:load COUNT10000000 # 向 bench-events 主题生产 1000 万条事件COUNT默认1000000。该任务Taskfile.yaml 的bench:load依次执行TRUNCATE bench_events—— 清空 MySQL 表保证计数从零开始删除并重建bench-events主题16 分区、复制因子 1用本仓库的 Redpanda Connect 二进制运行./producer.yaml造数go run ../../../../../../cmd/redpanda-connect/main.go run ./producer.yaml4.2 造数格式Kafka Connect 的 JSON Schema 信封这是整个基准最关键的细节。producer.yaml 生成的每条消息都不是裸 JSON而是Kafka Connect 的 schema/payload 信封格式目的是让 JDBC Sink 在没有 Schema Registry的情况下也能把字段映射到 MySQL 列{ schema: { type: struct, fields: [ {field: id, type: int64, optional: false}, {field: category, type: string, optional: true}, {field: value, type: double, optional: true}, {field: ts, type: int64, optional: true} ], optional: false, name: bench_events }, payload: { id: 1, category: cat_3, value: 8342.11, ts: 1736901234567890 } }信封由 Bloblangmapping生成id用count(events)递增、category用random_int(min: 1, max: 10)造 9 类、value用random_int(min: 1, max: 1000000).float64() / 100.0造浮点、ts用now().ts_unix_micro()打时间戳。输出端kafka_franz采用 snappy 压缩并以“3000 条或 4 MiB 或 5s”任一条件触发的批处理方式写出。4.3 为什么要用信封格式对照连接器配置见下一节中的value.converter: JsonConverter与value.converter.schemas.enable: true即可理解开启 schema 的 JsonConverter 会把信封里的schema部分解析为结构化类型信息从而把payload字段精确映射到bench_events的四列。这是不依赖 Confluent Schema Registry 时让 JDBC Sink 正常工作的标准做法。五、跑测注册连接器并测量写吞吐5.1 命令task bench:run TASKS16 # 注册 JDBC Sink 连接器并测量写吞吐TASKS设置连接器的tasks.max默认 16。bench:run任务自动完成一轮完整闭环删除旧连接器DELETE /connectors/mysql-jdbc-sink-bench并TRUNCATE bench_events让计数归零把消费组connect-mysql-jdbc-sink-bench的 offset 重置到 earliestkafka-consumer-groups --reset-offsets --to-earliest保证连接器从主题头完整重读用kafka-get-offsets读取主题末端 offset求和得到消息总数TOTAL通过 Kafka Connect REST API 注册连接器用jq现场改写tasks.maxcurl -sf -X POST http://localhost:8083/connectors \ -H Content-Type: application/json \ -d $(jq .config[\tasks.max\] \$TASKS\ connector.json) | jq -r .name从注册时刻起计时每 5sINTERVAL可覆盖查询一次 MySQL 行数打印TIME / WRITTEN / MSG/S表格当行数 TOTAL时输出汇总总消息数、总耗时、平均吞吐msg/s。5.2 连接器配置逐项解读connector.json 是基准的核心配置逐项说明如下配置项值含义与基准考量connector.classio.confluent.connect.jdbc.JdbcSinkConnector使用 Confluent JDBC Sink 连接器tasks.max16命令行可覆盖并行任务数直接决定写并行度是本次基准的自变量topicsbench-events消费源主题connection.urljdbc:mysql://mysql:3306/benchdbMySQL JDBC 连接串走 compose 内网connection.user/connection.passwordroot/password连接凭据insert.modeinsert只插入、不做 upsert避免额外开销pk.modenone不依赖主键纯追加写auto.createfalse表已由mysql-setup预建关闭自动建表table.name.formatbench_events目标表名value.converterorg.apache.kafka.connect.json.JsonConverter解析 JSON 信封value.converter.schemas.enabletrue读取信封内 schema 完成字段映射batch.size10000每批写入条数结果表中对比过 3000 与 10000consumer.override.fetch.min.bytes10485761 MiB抬高拉取门槛让轮询聚合更多数据consumer.override.fetch.max.wait.ms500拉取等待上限控制延迟与吞吐的折中consumer.override.max.poll.records5000单次 poll 最大记录数consumer.override.*前缀意味着这些参数会覆盖worker 全局的 consumer 配置只作用于该连接器是专为吞吐基准准备的调优点。5.3 辅助运维任务Taskfile 还提供若干辅助命令便于观察与排障task connector:create TASKS16 # 单独注册连接器不自动计时 task connector:status # 查看连接器与任务状态REST /connectors/.../status task connector:delete # 删除连接器不停止基础设施 task bench:mysql-count # 直接查询 bench_events 行数 task bench:measure INTERVAL5 # 独立版轮询测量复用同一套计数脚本 task logs:connect # 跟踪 Kafka Connect worker 日志 task logs:mysql # 跟踪 MySQL 日志六、组合跑测与收尾6.1 跑完整组合原文档给出的标准组合是依次测试 4 / 8 / 16 三个并行度task bench:load COUNT10000000 task bench:run TASKS4 task bench:run TASKS8 task bench:run TASKS16每次bench:run都会自动清理连接器、清空表、重置消费组 offset因此可以直接串行执行无需手动复位。6.2 停止与清理task down # 停止并移除所有容器与卷对应docker compose down -v --remove-orphans-v会一并删除数据卷确保下次task up从干净状态开始。七、结果怎么看与 Redpanda Connect 的横向对比完整结果见 docs/benchmark-results/mysql-cdc.md 的 “Kafka Connect JDBC Sink Comparison” 一节。该节记录了在 Intel Core i7-10850H 2.70GHz / 32 GB RAM / WSL2 环境下1000 万行从 Kafka 写入 MySQL 的实测吞吐batch.size 3000tasks.maxmsg/sec418,518831,2501642,553batch.size 10000tasks.maxmsg/sec1643,859结论要点引自结果文档属于该文档陈述的实测观察仅适用于该测试环境batch3000 时峰值42,553 msg/s16 tasks提升到 batch10000 只有边际收益43,859 msg/s说明批大小不是瓶颈任务数扩展收益递减明显4→8 tasks 约 1.7×8→16 tasks 约 1.4×作为对照同一份文档中 Redpanda Connectkafka_franz→sql_insert配置见 benchmark_config.yaml在 4 核、batch10000 时达到约64,102 msg/s约为 JDBC Sink 峰值的 1.5 倍且只用了 4 核而非 16 个 task更极端地Redpanda Connectmysql_cdc直读不经 Kafka峰值约191K msg/s约为 JDBC Sink 的 4.5 倍。需要强调的是以上数字均来自仓库内结果文档记录的既定测试环境不同机器、不同 MySQL 配置下绝对数值会变化但“JDBC Sink 的并行扩展先快后平、批大小非瓶颈、写库侧是最终瓶颈”这组结构性观察具备参考价值。八、从源码看 sql_insert对比基线另一侧的实现既然本基准的目的是为 Redpanda Connectsql_insert提供对比基线这里补充其底层实现要点出自 output_sql_insert.go每个消息执行一次行插入属于批量输出组件service.MustRegisterBatchOutput(sql_insert, ...)核心配置字段driver、dsn、table、columns、args_mapping、可选的prefix/suffix/options以及并行度max_in_flight默认 64见第 64-66 行的Field(service.NewIntField(max_in_flight)...Default(64))内部用Masterminds/squirrel构建INSERT语句args_mapping是一个 Bloblang 表达式求值结果必须是长度与columns一致的数组它支持 BatchPolicy 字段版本 3.59.0 起对应 rpcn 基准里batching: count: ${BATCH:10000}, period: 1s的用法见 benchmark_config.yaml这与 JDBC Sink 的batch.size参数形成天然对照。对照 JDBC Sink 的“N 个 task 各持有独立连接 batch.size 批插”模型sql_insert走的是“单进程内max_in_flight并发 批策略聚合”模型这正是两条路径资源占用与吞吐表现差异的根源也解释了为什么结果文档中后者能以更少核数取得更高吞吐。九、复现与扩展建议按以下步骤即可在本地完整复现本基准仓库只读以下均为运行而非修改# 1. 确保 Docker 运行进入基准目录 cd internal/impl/mysql/bench/mysql-write/jdbc-sink # 2. 起环境构建含 JDBC 插件的 Connect 镜像 task up # 3. 造数 1000 万 task bench:load COUNT10000000 # 4. 按不同并行度跑测 task bench:run TASKS4 task bench:run TASKS8 task bench:run TASKS16 # 5. 如需对比切到 rpcn 目录跑 Redpanda Connect 一侧 cd ../rpcn task up task bench:load COUNT10000000 task bench:run CORES4 BATCH10000 # 6. 清理 task down # 各自目录下执行扩展调优时建议优先观察tasks.max与max_in_flight的扩展曲线、batch.size/batching.count是否触及写库瓶颈、MySQL 侧innodb_buffer_pool_size与磁盘类型对绝对吞吐的影响。若想对照 CDC 读路径仓库还提供了 mysql-read 系列基准Debezium 与 rpcn 对比同一份 mysql-cdc.md 中均有记录与操作入口。十、小结Kafka Connect JDBC Sink 是 Kafka → MySQL 落库的事实标准组件之一而本基准的价值在于把它量化为可复现的基线固定数据模型、固定分区数、固定消息格式仅改变tasks.max与batch.size即可得到清晰的扩展曲线。结合 connector.json、Taskfile.yaml 与结果文档你可以随时在本仓库内重跑实验、验证结论并与sql_insert路径做直接对比为生产环境的写入链路选型提供实测依据。【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考