From 23796c901323bf823f06f858ca14ee04e4d40f69 Mon Sep 17 00:00:00 2001 From: anlicheng <244108715@qq.com> Date: Tue, 12 May 2026 16:53:41 +0800 Subject: [PATCH] fix outbox metadata --- apps/efka/src/iot/efka_iot_outbox.erl | 40 ++++++++++++++++++--------- 1 file changed, 27 insertions(+), 13 deletions(-) diff --git a/apps/efka/src/iot/efka_iot_outbox.erl b/apps/efka/src/iot/efka_iot_outbox.erl index 1b44b93..40a4b83 100644 --- a/apps/efka/src/iot/efka_iot_outbox.erl +++ b/apps/efka/src/iot/efka_iot_outbox.erl @@ -8,10 +8,10 @@ %%% %%% 文件结构: %%% - `segment-N.log':一个 segment 文件,按写入顺序追加保存 record。 -%%% - `metadata.term':保存 #{write_seq => N, acked_seq => N, -%%% writer_segment => N}。进程启动时仍然会扫描 segment 文件,并以 segment -%%% 文件里的真实数据作为准确信息;因此即使 metadata 落后,也不会导致跳过或 -%%% 删除未消费的数据。 +%%% - `metadata.json':使用原始 JSON 文本保存 write_seq、acked_seq 和 +%%% writer_segment。进程启动时仍然会扫描 segment 文件,并以 segment 文件 +%%% 里的真实数据作为准确信息;因此即使 metadata 落后,也不会导致跳过或删除 +%%% 未消费的数据。 %%% - outbox 目录通过 open/1 的 dir 参数传入,本模块不直接读取应用环境变量。 %%% %%% open/1 参数: @@ -89,7 +89,7 @@ -type outbox() :: #outbox{}. --define(METADATA_FILE, "metadata.term"). +-define(METADATA_FILE, "metadata.json"). -define(SEGMENT_PREFIX, "segment-"). -define(SEGMENT_EXT, ".log"). -define(RECORD_MAGIC, <<"EFKAOBX1">>). @@ -222,18 +222,32 @@ positive_option(Key, Options) -> read_metadata(MetadataPath) -> case file:read_file(MetadataPath) of {ok, Bin} -> - safe_metadata(binary_to_term(Bin, [safe])); + decode_metadata(Bin); {error, _} -> {0, 0, 1} end. +-spec decode_metadata(binary()) -> {non_neg_integer(), non_neg_integer(), pos_integer()}. +decode_metadata(Bin) -> + try json:decode(Bin) of + Metadata -> + safe_metadata(Metadata) + catch + _:_Reason -> + {0, 0, 1} + end. + -spec safe_metadata(term()) -> {non_neg_integer(), non_neg_integer(), pos_integer()}. -safe_metadata(#{write_seq := WriteSeq, acked_seq := AckedSeq, writer_segment := WriterSegment}) +safe_metadata(#{ + <<"write_seq">> := WriteSeq, + <<"acked_seq">> := AckedSeq, + <<"writer_segment">> := WriterSegment +}) when is_integer(WriteSeq), WriteSeq >= 0, is_integer(AckedSeq), AckedSeq >= 0, is_integer(WriterSegment), WriterSegment > 0 -> {WriteSeq, AckedSeq, WriterSegment}; -safe_metadata(#{write_seq := WriteSeq, acked_seq := AckedSeq}) +safe_metadata(#{<<"write_seq">> := WriteSeq, <<"acked_seq">> := AckedSeq}) when is_integer(WriteSeq), WriteSeq >= 0, is_integer(AckedSeq), AckedSeq >= 0 -> {WriteSeq, AckedSeq, 1}; safe_metadata(_) -> @@ -246,11 +260,11 @@ persist_metadata(#outbox{ write_seq = WriteSeq, acked_seq = AckedSeq }) -> - Metadata = term_to_binary(#{ - write_seq => WriteSeq, - acked_seq => AckedSeq, - writer_segment => WriterSegment - }), + Metadata = iolist_to_binary(json:encode(#{ + <<"write_seq">> => WriteSeq, + <<"acked_seq">> => AckedSeq, + <<"writer_segment">> => WriterSegment + })), TmpPath = MetadataPath ++ ".tmp", case file:write_file(TmpPath, Metadata, [write, binary]) of ok ->