Skip to content

Commit 2123136

Browse files
committed
焼く素材の行数を meta に載せ、途中で切れたら焼かせない
流し始めたあとに落ちても HTTP のステータスは変えられない(1 度しか送れない)ため、 途中で切れた素材は、受け取る側から見ると「短いだけの正しい素材」になる。 meta の min_docs を 1 で固定していたので 1 件でも検証を通り、そのまま焼けて 前の世代が捨てられていた。 本番でこれが起きた。68 万件を流した回が 1.7 万件で焼き上がり、残りの 66 万件が 消えている(次の回が読んだ前世代が 1.7 万件だったことと、区画の台帳に 68 万件を 割ったときの細かい升目が中身 0 で 328 個残っていることで確かめた)。 1 周目が数えた行数をそのまま下限として渡す。足りなければ取り込み側の検証 (世代の切り替えより前)で落ち、前の世代がそのまま残って記録にエラーが出る。 期限の境目も 1 周目のものを使い回す。2 周目で測り直すと、境目の 1 件が回を 跨いだだけで数が食い違い、正しく焼けたものまで弾かれる。 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1 parent c86946b commit 2123136

3 files changed

Lines changed: 37 additions & 4 deletions

File tree

app/collect.py

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3314,7 +3314,17 @@ def bake_survey(item, sources: dict, previous, collected, only_new=False, edits=
33143314
src = sources.get(rule["source"])
33153315
if src is not None and wanted[n]:
33163316
alive[n] = _existing_titles(src.path, wanted[n])
3317-
plan = {"diff": diff, "rules": rules, "alive": alive, "first_title": first_title}
3317+
plan = {
3318+
"diff": diff, "rules": rules, "alive": alive, "first_title": first_title,
3319+
# **2 周目に流す行数**。取り込み側の検証の下限になる(`bake_lines` の meta)——
3320+
# 流し始めたらステータスは変えられないので、途中で切れた素材と最後まで届いた
3321+
# 素材は、受け取った側からは見分けが付かない。数だけが手掛かりになる。
3322+
# 数えるのは期限で落としたあと(流すのもそのあと)
3323+
"rows": total,
3324+
# **期限の境目は 1 周目のものを使い回す。** 2 周目で測り直すと、境目の
3325+
# 1 件が回を跨いだだけで数が食い違い、正しく焼けたものが弾かれる
3326+
"limit": limit,
3327+
}
33183328
# **落としたタグの数も、流し始める前に数える。** 控えに残すのはここで record する
33193329
# ためで、流しながら数えると「控えを書いたあとに分かる」ことになる。
33203330
# **指定を持つ収集だけ**もう 1 周する(持たない収集では 1 件も落ちない)
@@ -3353,12 +3363,18 @@ def bake_lines(item, sources: dict, previous, collected, only_new=False, edits=F
33533363
"""
33543364
plan = survey or bake_survey(item, sources, previous, collected, only_new, edits)
33553365
rows = _rows_of(previous)
3356-
limit = _expiry_limit(item)
3366+
limit = plan.get("limit", _expiry_limit(item))
33573367

33583368
yield json.dumps({
33593369
"meta": {
33603370
"dump_date": _dump_date(item.name, sources),
3361-
"min_docs": 1,
3371+
# **これから流す行数をそのまま下限にする。** 流し始めたあとに落ちても
3372+
# ステータスは変えられない(1 度しか送れない)ので、**途中で切れた素材は
3373+
# 受け取る側から見ると「短いだけの正しい素材」**になる —— 焼けてしまうと
3374+
# 前の世代は捨てられ、届かなかったぶんは消える。
3375+
# 本番でこれが起きた(68 万件のうち 1.7 万件で焼き上がり、残りが消えた)。
3376+
# 数が足りなければ取り込み側が検証で落とし、前の世代がそのまま残る
3377+
"min_docs": max(1, int(plan.get("rows") or 1)),
33623378
"sample_titles": [plan["first_title"]],
33633379
}
33643380
}, ensure_ascii=False)

docs/adding-a-source.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,6 +184,11 @@ GET /fetch?source=NAME → NDJSON(1 行目に meta、以降は 1 行 1 文書)
184184
```
185185

186186
- 文書の項目はコアスキーマの `Doc` と同じ(`title` だけ必須)。`doc_id` を省くと行番号を振ります
187+
- **`min_docs` には「これから流す行数」を入れてください。** 流し始めたあとに落ちても
188+
HTTP のステータスは変えられないので、**途中で切れた素材は、受け取る側から見ると
189+
「短いだけの正しい素材」**になります —— そのまま焼けると前の世代は捨てられ、
190+
届かなかったぶんは消えます(集める層でこれが起き、68 万件のうち 1.7 万件で
191+
焼き上がりました)。数が足りなければ構築後の検証で落ち、前の世代がそのまま残ります
187192
- **`meta` をヘッダではなく本文の 1 行目に置く**のは、HTTP ヘッダが latin-1 しか運べず
188193
日本語のタイトルを載せられないため。取り込んだ中身を見てから代表を選ぶソースは、
189194
全部数え終えてから 1 行目を書き出せます

tests/test_collect.py

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4759,9 +4759,21 @@ def test_it_yields_one_line_at_a_time(self):
47594759
assert not isinstance(lines, list)
47604760
out = list(lines)
47614761
# 1 行目は meta、以降が 1 行 1 文書
4762-
assert json.loads(out[0])["meta"]["min_docs"] == 1
4762+
assert json.loads(out[0])["meta"]["min_docs"] == 3
47634763
assert [json.loads(line)["title"] for line in out[1:]] == ["店1", "店2", "店3"]
47644764

4765+
def test_the_meta_says_how_many_lines_follow(self):
4766+
"""**途中で切れた素材を、短いだけの正しい素材として焼かせない。**
4767+
4768+
流し始めたあとに落ちてもステータスは変えられない(1 度しか送れない)ので、
4769+
受け取る側から見ると区別が付かない —— 焼けてしまうと前の世代は捨てられ、
4770+
届かなかったぶんは消える。本番でこれが起きた(68 万件のうち 1.7 万件で
4771+
焼き上がり、残りが消えた)。数だけが手掛かりになる。
4772+
"""
4773+
out = list(collect.bake_lines(self.item(), {}, self.rows(9), []))
4774+
4775+
assert json.loads(out[0])["meta"]["min_docs"] == len(out) - 1
4776+
47654777
def test_the_same_shape_as_the_whole_string(self):
47664778
# 丸ごと組む道(`ndjson`)と、流す道が同じものを返すこと
47674779
whole, _diff = collect.ndjson(self.item(), {}, self.rows(3), [])

0 commit comments

Comments
 (0)