设计与实现Design
代码不多,边界很清楚:Kafka 的事交给 moonkafka。
There is not much code, and the boundary is sharp: Kafka belongs to moonkafka.
模块边界Module boundary
cmd/main daqing/moonkafka Apache Kafka
──────────────────────────── ────────────────────────────── ────────────
argument parsing broker connection broker
UTF-8 conversion ────────► metadata handling ────────► topic log
result rendering Kafka protocol
Producer::send
Consumer::poll
左边这一列就是本项目全部的职责。协议编码、连接管理、元数据刷新、偏移量跟踪都在
moonkafka 里,cmd/main 没有一行自定义的 Kafka 客户端逻辑。
The left column is the whole of this project's responsibility. Protocol encoding,
connection management, metadata refresh and offset tracking all live in moonkafka;
cmd/main contains no custom Kafka client logic at all.
| 包 | 导入 |
|---|---|
cmd/main | daqing/moonkafka、moonbitlang/async、core/encoding/utf8、core/env、core/string |
tests/itest | async、async/fs、async/shell、core/argparse、core/string |
| Package | Imports |
|---|---|
cmd/main | daqing/moonkafka, moonbitlang/async, core/encoding/utf8, core/env, core/string |
tests/itest | async, async/fs, async/shell, core/argparse, core/string |
模块依赖是 daqing/moonkafka@0.3.4 与
moonbitlang/async@0.22.4。moon.mod 声明
preferred_target = "native";两个包都限制
supported_targets = "+native",cmd/main 还为 native 目标加了
cc-link-flags = "-lz"。
The module depends on daqing/moonkafka@0.3.4 and
moonbitlang/async@0.22.4. moon.mod declares
preferred_target = "native"; both packages constrain
supported_targets = "+native", and cmd/main adds
cc-link-flags = "-lz" for the native target.
produce 的路径The produce path
取参数、建连接、发一条、打印偏移量——全部发生在
@async.with_task_group 里,连接在作用域结束时由 defer 关闭。
Read the arguments, connect, send once, print the offset — all inside an
@async.with_task_group, with the connection closed by
defer when the scope ends.
async fn run_producer(args : Array[String]) -> Unit {
guard args.get(2) is Some(topic) && args.get(3) is Some(value) else {
println("usage: produce <topic> <value> [key] [host] [port]")
return
}
let key = args.get(4).map(text => @utf8.encode(text))
let host = args.get(5).unwrap_or("127.0.0.1")
let port = parse_port(args, 6)
@async.with_task_group(group => {
let producer = @moonkafka.Producer::connect(
group~, host~, port~, topic~, acks=-1,
)
defer producer.close()
let offset = producer.send(key?, value=@utf8.encode(value))
println("sent to topic \"\{topic}\" at offset \{offset}")
})
}
key是可选的:省略时key?传None,发送一条无 key 的记录。acks=-1要求所有同步副本确认,Producer::send因此返回 broker 分配的偏移量(Int64),直接打印出来。- key 与 value 都先经
@utf8.encode转成Bytes——moonkafka 的 API 收发的是字节。
keyis optional: when omitted,key?isNoneand a keyless record is sent.acks=-1requires every in-sync replica to acknowledge, soProducer::sendreturns the offset the broker assigned (Int64), which is printed directly.- Both key and value go through
@utf8.encodefirst — moonkafka's API takes and returns bytes.
consume 的路径The consume path
async fn run_consumer(args : Array[String]) -> Unit {
let (max_messages, args) = take_max_messages(args)
guard args.get(2) is Some(topic) else { /* usage */ return }
let host = args.get(3).unwrap_or("127.0.0.1")
let port = parse_port(args, 4)
let start_from = parse_start_from(args.get(5))
@async.with_task_group(group => {
let consumer = @moonkafka.Consumer::connect(
group~, host~, port~, topic~, start_from~,
)
defer consumer.close()
let mut received = 0
for ;; {
for record in consumer.poll() {
let key = record.key_utf8().unwrap_or("<null>")
let value = record.value_utf8().unwrap_or("<tombstone>")
println("offset=\{record.offset} timestamp=\{record.timestamp} ...")
received = received + 1
}
if max_messages is Some(limit) && received >= limit {
break
}
}
})
}
外层是无限循环,内层遍历 poll() 返回的一批记录。
--max-messages 的判定放在批处理循环之外:计数达标就跳出,
因此读够 n 条后退出码是 0,不会因为多读了一批而超发。
The outer loop is unbounded; the inner one walks the batch poll()
returns. The --max-messages check sits outside the batch loop:
once the count is reached the program breaks, so it exits 0 after exactly n records
rather than overshooting by a batch.
这是一个「简单消费者」This is a simple consumer
moonkafka 的 Consumer 在注释里写得很清楚:它拉取单个主题的所有分区,
不做 consumer group 协调,偏移量只记在内存里。这对一个演示程序正合适,
但有两件事需要知道:
moonkafka's Consumer says so itself: it fetches all partitions of one
topic, without consumer group coordination, and tracks offsets in
memory. That suits a demo, with two consequences worth knowing:
- 偏移量不会被提交,所以重启后不会从上次的位置继续——一切都由
earliest/latest决定。 - 多个实例之间不会分摊分区,每个实例都会读到自己那份全量数据。
- Offsets are never committed, so a restart does not resume where it left off —
earliest/latestdecides everything. - Several instances do not split partitions between them; each one reads its own copy of everything.
Consumer::poll() 接受 max_wait_ms(默认 500)与
max_bytes(默认 1 MiB)两个可选参数,本 demo 都用默认值。
Consumer::poll() takes optional max_wait_ms (500 by default)
and max_bytes (1 MiB by default); the demo uses both defaults.
参数解析Argument parsing
--max-messages 需要前置处理:它在解析早期被
take_max_messages 摘出来,返回「限制值 + 剩余位置参数」。
--max-messages needs a pre-pass:
take_max_messages pulls it out early and returns the limit together with
the remaining positional arguments, in order.
fn take_max_messages(args : Array[String]) -> (Int?, Array[String])
- 同时接受
--max-messages 5与--max-messages=5。 - 非 flag 的参数按原顺序保留,因此索引位置不变(
args[2]仍是 topic)。 - 缺失取值或取值非法都会当场中止。
- Both
--max-messages 5and--max-messages=5are accepted. - Non-flag arguments keep their original order, so indices do not shift (
args[2]is still the topic). - A missing or invalid value aborts immediately.
--max-messages
placed in the middle shifts later positionals up by one. Details in the
CLI reference.
parse_port 与 parse_start_from 在同一层做校验:端口用
@string.parse_int 转换并在 catch 里报错,起始位置用模式匹配
限定为两个取值。
parse_port and parse_start_from validate at the same layer:
the port is converted with @string.parse_int and reports errors from a
catch block, while the start position is pattern-matched down to its two
legal values.
错误处理Error handling
| 场景 | 做法 |
|---|---|
| 缺少必需参数 | guard … else { println(usage); return }——打印用法后正常结束 |
| 取值非法 | println(…) ; abort("")——先给出人能看懂的说明,再中止 |
| 数字解析失败 | @string.parse_int(v) catch { _ => … },把异常换成一条消息 |
| 资源释放 | defer consumer.close() / defer producer.close(),连接一定被关掉 |
| 任务生命周期 | @async.with_task_group(group => …)——moonkafka 的连接需要借一个 task group |
| Situation | Approach |
|---|---|
| Missing required argument | guard … else { println(usage); return } — print usage, return normally |
| Invalid value | println(…) ; abort("") — a readable message first, then abort |
| A number fails to parse | @string.parse_int(v) catch { _ => … }, turning the exception into a message |
| Releasing resources | defer consumer.close() / defer producer.close() — the connection is always closed |
| Task lifetime | @async.with_task_group(group => …) — moonkafka's connections borrow a task group |
harness 的设计How the harness is built
命令是值,不是副作用Commands are values, not side effects
一条外部命令建模为 Step:用途、程序、参数。它由一批共享的构造函数生成,
--dry-run 只是把这些值渲染出来,真实运行才去执行它们。
因为两边同源,计划不可能与真实行为脱节。
An external command is modelled as a Step: a purpose, a program, and its
arguments. Shared builders produce those values; --dry-run merely
renders them while a real run executes them. Because both sides come
from the same source, the plan cannot drift from what actually happens.
struct Step {
purpose : String
program : String
args : Array[String]
}
三种执行语义Three execution semantics
| 函数 | 用途 |
|---|---|
run_required | harness 依赖的前置步骤。失败就中止整个运行——后续阶段的结果没有意义 |
run_observed | 结果本身就是断言的一部分。无论成败都记录退出码与耗时 |
run_until_success | 对 produce 的有界重试,用来吸收 broker 的启动竞态 |
| Function | Purpose |
|---|---|
run_required | A step the harness depends on. Failure aborts the whole run — later phases could not mean anything |
run_observed | The outcome is itself part of the assertion. Exit code and duration are recorded either way |
run_until_success | Bounded retries around the produce, absorbing the broker's startup race |
断言与中止Assertions and aborting
Runner 只记三件事:跑了多少条断言、失败多少、是否已经中止。
check 记录一条结论;fail 额外把运行标记为中止,后续阶段直接跳过。
Runner tracks three things: how many checks ran, how many failed, and
whether the run has already aborted. check records a verdict;
fail additionally marks the run aborted so later phases are skipped.
进程的失败通过 suberror ItestError::Aborted(String) 表达,并且刻意放在
main 的最后一行抛出——那时容器已经停掉、资源已经释放,
进程才带着一行可读的消息以非零码退出。
Failure is expressed as suberror ItestError::Aborted(String), raised
deliberately on the last line of main: by then the container is
stopped and resources are released, and the process exits non-zero with one readable
message.
外部命令不会挂住 harnessExternal commands cannot hang the harness
每条命令都经 @shell.Cmd(...).output(timeout_ms~) 执行,带硬超时。
连「程序根本起不来」这种情况也算作一次普通的失败步骤(exit_code = -1 加一条说明),
而不是抛异常——比如一个没装的 compose provider。
Every command runs through @shell.Cmd(...).output(timeout_ms~) under a
hard timeout. Even "the program could not be spawned at all" counts as an ordinary
failed step (exit_code = -1 plus an explanation) rather than an
exception — a compose provider that is not installed, say.