É-01: Ţĥé þŕöçéššîñĝ éñĝîñé
Šüḿḿàŕý
Çöñţéñţ îš þŕöçéššéđ ƃý à çĥàññéļ-ƃàšéđ šţŕéàḿîñĝ þîþéļîñé. Éàçĥ ţööļ ŕüñš îñ
îţš öŵñ ĝöŕöüţîñé; ţööļš àŕé çöññéçţéđ ƃý ƃüƒƒéŕéđ çĥàññéļš ţĥàţ þŕöṽîđé
àüţöḿàţîç ƃàçķþŕéššüŕé. Àñ errgroup.Group çööŕđîñàţéš éŕŕöŕš àñđ þŕöþàĝàţéš
çöñţéẋţ çàñçéļļàţîöñ. Ƒöüŕ îñđéþéñđéñţ çöñçüŕŕéñçý ļàýéŕš (îñţŕà-ţööļ ƃļöçķ
þàŕàļļéļîšḿ, ƃàţçĥ ƒîļé çöñçüŕŕéñçý, đöçüḿéñţ-ļéṽéļ çöñçüŕŕéñçý, àñđ šţŕéàḿîñĝ
öƃšéŕṽàţîöñ) çöḿþöšé ŵîţĥöüţ îñţéŕƒéŕéñçé. Ƒļöŵš àŕé đéçļàŕéđ àš éîţĥéŕ à
ĝŕàþĥ öƒ ñöđéš àñđ éđĝéš öŕ à šéǫüéñţîàļ ļîšţ öƒ šţéþš (ŵîţĥ éẋþļîçîţ
parallel: ƃļöçķš ƒöŕ ƒàñ-öüţ), ƃöţĥ çöḿþîļéđ ţö ţĥé šàḿé éẋéçüţàƃļé
ŕéþŕéšéñţàţîöñ.
Çöñţéẋţ
Ĝö'š ĝöŕöüţîñéš àñđ çĥàññéļš ḿàķé îţ ñàţüŕàļ ţö šţŕüçţüŕé à þîþéļîñé àš çöñçüŕŕéñţ šţàĝéš çöññéçţéđ ƃý ţýþéđ çĥàññéļš. Šüçĥ à þîþéļîñé ḿîẋéš ÇÞÜ-ƃöüñđ ŵöŕķ (ƒöŕḿàţ þàŕšîñĝ, çĥéçķš) ŵîţĥ ÎÖ-ƃöüñđ ŵöŕķ (ĻĻḾ çàļļš, ḿéḿöŕý ļööķüþš). Ţĥé šàḿé þîþéļîñé ḿüšţ àļšö ƃé đŕîṽéñ àţ šéṽéŕàļ šçàļéš: à šîñĝļé ƒîļé öñ à ļàþţöþ; ĥüñđŕéđš öƒ ƒîļéš îñ à ƃàţçĥ; à ļöñĝ-ļîṽéđ þŕöĵéçţ ŵîţĥ ḿàñý đöçüḿéñţš þŕöçéššéđ îñ þàŕàļļéļ.
Îţ ḿüšţ àļšö šüþþöŕţ ƃöţĥ đéçļàŕàţîṽé àüţĥöŕîñĝ (à ṽîšüàļ ƒļöŵ éđîţöŕ, ĥüḿàñ-ŕéàđàƃļé ÝÀḾĻ) àñđ þŕöĝŕàḿḿàţîç çöñšţŕüçţîöñ (ţĥé Ĝö ļîƃŕàŕý) ƒŕöḿ öñé đàţà ḿöđéļ.
Đéçîšîöñ
Çĥàññéļ-ƃàšéđ þîþéļîñé
Çöñţéñţ ƒļöŵš ţĥŕöüĝĥ à çĥàññéļ-ƃàšéđ çöñçüŕŕéñţ þîþéļîñé:
Çöñţéñţ éñţéŕš ţĥŕöüĝĥ à šöüŕçé ƃîñđîñĝ àñđ ļéàṽéš ţĥŕöüĝĥ à šîñķ
ƃîñđîñĝ (É-04). Ƒöŕ ţĥé đéƒàüļţ file ƃîñđîñĝ
ţĥéšé àŕé à ĐàţàƑöŕḿàţ ŕéàđéŕ àñđ ŵŕîţéŕ (É-02); à
þŕöĵéçţ šţöŕé, à .kpz, öŕ àñ îñţéŕçĥàñĝé ƒîļé ƃîñđ ţĥé šàḿé šţŕéàḿ ŵîţĥ ñö
ŕéàđéŕ öŕ ŵŕîţéŕ. Ƃéţŵééñ ţĥé éñđš, éàçĥ ţööļ ŕüñš îñ îţš öŵñ ĝöŕöüţîñé.
Ƃüƒƒéŕéđ çĥàññéļš þŕöṽîđé ƃàçķþŕéššüŕé. errgroup.Group çööŕđîñàţéš éŕŕöŕ
ĥàñđļîñĝ àçŕöšš ĝöŕöüţîñéš, àñđ çöñţéẋţ çàñçéļļàţîöñ þŕöþàĝàţéš ţö àļļ šţàĝéš.
Þàŕţš çàŕŕý ţýþéđ ŕéšöüŕçéš (Ƒ-02): Ƃļöçķš çöñţàîñ ţŕàñšļàţàƃļé çöñţéñţ, Đàţà çàŕŕîéš šţŕüçţüŕàļ ḿàŕķüþ, Ļàýéŕš ĝŕöüþ ñéšţéđ çöñţéñţ, Ḿéđîà ĥöļđš ƃîñàŕý àššéţš. Ţööļš đéçļàŕé ŵĥîçĥ ŕéšöüŕçé ţýþéš ţĥéý ĥàñđļé; ţĥé ŕéšţ þàšš ţĥŕöüĝĥ üñçĥàñĝéđ.
Éẋéçüţöŕ
flow.DefaultExecutor öŕçĥéšţŕàţéš ţööļ çĥàîñš üšîñĝ ţĥé ĝöŕöüţîñé-þéŕ-ţööļ
ḿöđéļ:
- Éàçĥ ţööļ îñ ţĥé çĥàîñ ŕüñš îñ îţš öŵñ ĝöŕöüţîñé.
- Ƃüƒƒéŕéđ çĥàññéļš (çöñƒîĝüŕàƃļé šîžé) çöññéçţ àđĵàçéñţ ţööļš.
errgroupçöļļéçţš ţĥé ƒîŕšţ éŕŕöŕ àñđ çàñçéļš ţĥé šĥàŕéđ çöñţéẋţ.- Þàŕàļļéļ đöçüḿéñţ þŕöçéššîñĝ îš ƃöüñđéđ ƃý à šéḿàþĥöŕé ŵöŕķéŕ þööļ.
ToolFactoriesçŕéàţé ƒŕéšĥ ţööļ îñšţàñçéš þéŕ đöçüḿéñţ, šö çöñçüŕŕéñţ đöçüḿéñţš ñéṽéŕ šĥàŕé ḿüţàƃļé ţööļ šţàţé.
Çöñƒîĝüŕàţîöñ üšéš ţĥé ƒüñçţîöñàļ-öþţîöñš þàţţéŕñ:
executor := flow.NewExecutor(
flow.WithMaxConcurrency(8),
flow.WithChannelSize(128),
flow.WithCollectors(wordCounter, checkReport),
)
flow.WithImmutabilityCheck ţüŕñš öñ ţĥé þéŕ-đîšþàţçĥ çöñţéñţ-îḿḿüţàƃîļîţý
ƃàçķšţöþ ƒöŕ à ŕüñ (É-03); îţ îš đéṽ/ţéšţ ţööļîñĝ àñđ îš
öƒƒ üñļéšš à çàļļéŕ àšķš ƒöŕ îţ.
Þàŕàļļéļ ƃļöçķ þŕöçéššîñĝ
Ƒöŕ ÎÖ-ƃöüñđ ţööļš (ĻĻḾ ţŕàñšļàţîöñ, ŕéḿöţé çĥéçķš), šéǫüéñţîàļ þéŕ-þàŕţ
þŕöçéššîñĝ üñđéŕüţîļîžéš ţĥŕöüĝĥþüţ. tool.ParallelBlockTool ŵŕàþš àñý ţööļ ţö
ƒàñ öüţ Ƃļöçķ þŕöçéššîñĝ àçŕöšš Ñ ĝöŕöüţîñéš ŵĥîļé þŕéšéŕṽîñĝ šţŕîçţ Þàŕţ
öŕđéŕîñĝ:
Ţĥé đîšþàţçĥéŕ àššîĝñš ḿöñöţöñîç šéǫüéñçé ñüḿƃéŕš ţö àļļ îñçöḿîñĝ Þàŕţš. Ƃļöçķ Þàŕţš àŕé đîšþàţçĥéđ ţö à šéḿàþĥöŕé-ƃöüñđéđ ŵöŕķéŕ þööļ; ñöñ-Ƃļöçķ Þàŕţš (Đàţà, Ḿéđîà, Ļàýéŕ) þàšš ţĥŕöüĝĥ ţĥé îññéŕ ţööļ šéǫüéñţîàļļý. À ḿîñ-ĥéàþ ŕéàššéḿƃļý ƃüƒƒéŕ çöļļéçţš ŕéšüļţš àñđ éḿîţš ţĥéḿ îñ šţŕîçţ šéǫüéñçé öŕđéŕ, šö đöŵñšţŕéàḿ ţööļš šéé ţĥé šàḿé Þàŕţ öŕđéŕîñĝ ŕéĝàŕđļéšš öƒ ŵĥîçĥ ŵöŕķéŕ ƒîñîšĥéđ ƒîŕšţ.
Àüţö-þàŕàļļéļîšḿ îš à ţööļ þŕöþéŕţý. Éàçĥ ţööļ đéçļàŕéš
ToolMeta.DefaultParallelBlocks (É-03); ţĥé ŕüññéŕ ţàķéš
ţĥé ḿàẋîḿüḿ àçŕöšš ţĥé ƒļöŵ'š ţööļš àñđ ŵŕàþš éṽéŕý ţööļ àţ ţĥàţ ŵîđţĥ. À
þŕöĵéçţ ḿàý þîñ îţš öŵñ ṽàļüé, àñđ --parallel-blocks N öṽéŕŕîđéš ƃöţĥ
(--parallel-blocks 1 đîšàƃļéš ţĥé ŵŕàþþéŕ). Ţĥé ţööļš ţĥàţ đéçļàŕé à đéƒàüļţ
ţöđàý àŕé ţĥé ĻĻḾ-ƃàçķéđ öñéš, šö àñ öŕđîñàŕý ŕüļéš-öñļý ƒļöŵ ŕüñš šéǫüéñţîàļļý
ŵîţĥ ñö çöñƒîĝüŕàţîöñ.
Ƃàţçĥ éẋéçüţöŕ
flow.BatchExecutor þŕöçéššéš ḿüļţîþļé þŕé-ŕéàđ ƒîļéš ţĥŕöüĝĥ à ţööļ çĥàîñ ŵîţĥ
çöñƒîĝüŕàƃļé ƒîļé-ļéṽéļ çöñçüŕŕéñçý:
type BatchConfig struct {
FileConcurrency int // max files processed in parallel (default: 1)
ChannelSize int // per-pipeline channel buffer size (default: 64)
SharedResources []io.Closer // resources shared across files (closed at end)
FailFast bool // cancel remaining on first error (default: true)
}
Éàçĥ ƒîļé ĝéţš ƒŕéšĥ ţööļ îñšţàñçéš ƒŕöḿ ŢööļƑàçţöŕý ƒüñçţîöñš, þŕéṽéñţîñĝ šţàţé ļéàķàĝé ƃéţŵééñ çöñçüŕŕéñţ đöçüḿéñţš. Ŕéšüļţš àŕé ŕéţüŕñéđ îñ îñþüţ ƒîļé öŕđéŕ ŕéĝàŕđļéšš öƒ çöḿþļéţîöñ öŕđéŕ. Çöļļéçţöŕš àŕé çàļļéđ ŵîţĥ ḿüţéẋ þŕöţéçţîöñ ƒöŕ ţĥŕéàđ-šàƒé àĝĝŕéĝàţîöñ àçŕöšš ƒîļéš.
Çöñçüŕŕéñçý ļàýéŕîñĝ
Ƒöüŕ îñđéþéñđéñţ çöñçüŕŕéñçý ļàýéŕš çöḿþöšé ŵîţĥöüţ îñţéŕƒéŕéñçé:
| Ļàýéŕ | Šçöþé | Çöñţŕöļ | Öŕđéŕ |
|---|---|---|---|
| ÞàŕàļļéļƂļöçķŢööļ | Ƃļöçķš ŵîţĥîñ öñé ţööļ | Ñ ĝöŕöüţîñéš þéŕ ţööļ | Šţŕîçţ Þàŕţ öŕđéŕ |
| ƂàţçĥÉẋéçüţöŕ | Ḿüļţîþļé ƒîļéš | ƑîļéÇöñçüŕŕéñçý šéḿàþĥöŕé | Ƒîļé öŕđéŕ þŕéšéŕṽéđ |
| Éẋéçüţöŕ | Ḿüļţîþļé đöçüḿéñţš | ḾàẋÇöñçüŕŕéñçý šéḿàþĥöŕé | Đöçüḿéñţ öŕđéŕ þŕéšéŕṽéđ |
| ŢàþþîñĝŢööļ | Öƃšéŕṽàţîöñ | Îñļîñé (ñö éẋţŕà ĝöŕöüţîñé) | Šéǫüéñţîàļ |
Çöļļéçţöŕš àñđ šţŕéàḿîñĝ çöļļéçţöŕš
Çöļļéçţöŕš àĝĝŕéĝàţé ŕéšüļţš àçŕöšš đöçüḿéñţš (ŵöŕđ çöüñţš, çĥéçķ ŕéþöŕţš, ţéŕḿ
ļîšţš). Ţĥéý îḿþļéḿéñţ Collect(ctx, item, parts), çàļļéđ àƒţéŕ éàçĥ đöçüḿéñţ
çöḿþļéţéš, àñđ Result() ƒöŕ ţĥé ƒîñàļ àĝĝŕéĝàţé. Çöļļéçţöŕš ḿüšţ ƃé
ţĥŕéàđ-šàƒé šîñçé ḿüļţîþļé đöçüḿéñţš ḿàý çöḿþļéţé çöñçüŕŕéñţļý.
StreamingCollector éẋţéñđš Collector ŵîţĥ Observe(part) ƒöŕ îñļîñé
öƃšéŕṽàţîöñ ŵîţĥöüţ àđđîñĝ à þîþéļîñé šţàĝé. TappingTool ŵŕàþš à ţööļ àñđ îţš
šţŕéàḿîñĝ çöļļéçţöŕ: öüţþüţ Þàŕţš àŕé îñţéŕçéþţéđ àñđ þàššéđ ţö Observe()
šýñçĥŕöñöüšļý ƃéƒöŕé ƒöŕŵàŕđîñĝ đöŵñšţŕéàḿ. Ţĥîš éñàƃļéš ŕéàļ-ţîḿé ḿéţŕîçš
ŵîţĥöüţ ƃüƒƒéŕîñĝ ţĥé éñţîŕé ŕéšüļţ šéţ.
Ƒļöŵ ţŕàçîñĝ àñđ ṽîšüàļîžàţîöñ
flow.TraceRecorder çàþţüŕéš ţîḿéšţàḿþéđ éṽéñţš đüŕîñĝ ƒļöŵ éẋéçüţîöñ.
flow.TracingTool ŵŕàþš éàçĥ ţööļ îñ ţĥé çĥàîñ àñđ ŕéçöŕđš éñţéŕ/éẋîţ éṽéñţš
ŵîţĥ Þàŕţ šñàþšĥöţš. Ţĥé --trace path/to/trace.json ƒļàĝ öñ kapi run éñàƃļéš
ţŕàçîñĝ. Ţĥé öüţþüţ îš à FlowTrace ĴŠÖÑ ƒîļé çöñţàîñîñĝ:
- Ñöđéš: ţĥé ţööļ çĥàîñ ŵîţĥ çöñçüŕŕéñçý ḿéţàđàţà
- Éṽéñţš: ţîḿéšţàḿþéđ éñţéŕ/éẋîţ éṽéñţš þéŕ Þàŕţ
- Þàŕţ šñàþšĥöţš: Þàŕţ šţàţé ƃéƒöŕé àñđ àƒţéŕ éàçĥ ñöđé
- Đüŕàţîöñ: ţöţàļ ƒļöŵ éẋéçüţîöñ ţîḿé îñ ḿîçŕöšéçöñđš
À ƃŕöŵšéŕ-ƃàšéđ ṽîšüàļîžàţîöñ ŕéñđéŕš ţĥé ţŕàçé àš àñ àñîḿàţéđ ƒļöŵ đîàĝŕàḿ ŵîţĥ þàŕţîçļéš ḿöṽîñĝ ţĥŕöüĝĥ ñöđéš, çĥàññéļ ƒîļļ îñđîçàţöŕš, àñđ ŵöŕķéŕ ļàñé šéþàŕàţîöñ ƒöŕ þàŕàļļéļ ţööļš. Ţĥé þļàýƃàçķ éñĝîñé šüþþöŕţš ṽàŕîàƃļé-šþééđ ŕéþļàý àñđ šééķîñĝ.
Öƃšéŕṽàţîöñ šéàḿ
Ţŕàçîñĝ ţö à ƒîļé îš öñé çöñšüḿéŕ öƒ à ḿöŕé ĝéñéŕàļ šéàḿ. core/observe
đéƒîñéš à Tracer / Span îñţéŕƒàçé ŵîţĥ à ñö-öþ đéƒàüļţ àñđ ñö đéþéñđéñçý öñ
àñý ţéļéḿéţŕý ļîƃŕàŕý, šö îţ çöšţš ñöţĥîñĝ üñţîļ à ĥöšţ ŕéĝîšţéŕš à ţŕàçéŕ àţ
šţàŕţüþ; kapi àñđ ţĥé đéšķţöþ ŕéĝîšţéŕ ñöñé. flow.WrapWithSpans ŵŕàþš éàçĥ
ţööļ îñ à çĥàîñ šö ţĥàţ öñé šþàñ öþéñš þéŕ ţööļ îñṽöçàţîöñ (flow.tool), àñđ
ţĥé ƒîļé ŕüññéŕ öþéñš öñé þéŕ ƒöŕḿàţ ŕéàđ àñđ ŵŕîţé (format.read,
format.write). À SessionTool ķééþš îţš šéššîöñ þàţĥ ţĥŕöüĝĥ ţĥé ŵŕàþþéŕ, šö
à ŕéšüḿàƃļé ŕüñ šţîļļ ŕéšüḿéš. Ḿöđéļ þŕöṽîđéŕ çàļļš çàŕŕý ţĥéîŕ öŵñ šþàñ ţĥŕöüĝĥ
ţĥé šàḿé šéàḿ (É-07). Öñé šþàñ þéŕ ţööļ ŕàţĥéŕ ţĥàñ
þéŕ Þàŕţ ķééþš ţĥé çöšţ þŕöþöŕţîöñàļ ţö ţĥé ĥàñđƒüļ öƒ ţööļš à ƒļöŵ ĥàš, ñöţ ţĥé
ţĥöüšàñđš öƒ Þàŕţš îţ ḿöṽéš.
Ƒļöŵ đéƒîñîţîöñš
flow.FlowDefinition îš à ĴŠÖÑ/ÝÀḾĻ-šéŕîàļîžàƃļé šţŕüçţ ţĥàţ çàþţüŕéš à ƒļöŵ
ĝŕàþĥ (ñöđéš + éđĝéš) àñđ ţĥé ţööļ çöñƒîĝüŕàţîöñš ñééđéđ ţö ŕéçöñšţŕüçţ à
ŕüññàƃļé ƒļöŵ. Ţĥîš šéþàŕàţéš ţĥé đéçļàŕàţîṽé đéšçŕîþţîöñ öƒ à ƒļöŵ ƒŕöḿ îţš
ŕüñţîḿé éẋéçüţîöñ.
Éàçĥ FlowNode ĥàš:
- ÎĐ: üñîǫüé îđéñţîƒîéŕ ŵîţĥîñ ţĥé đéƒîñîţîöñ
- Ţýþé:
flow.NodeToolƒöŕ à þŕöçéššîñĝ šţéþ.NodeReaderàñđNodeWriterŕéḿàîñ îñ ţĥéNodeTypeṽöçàƃüļàŕý ƒöŕ ţĥé éđîţöŕ'š ĝŕàþĥ, ƃüţ éṽéŕý ƃüîļţ-îñ àñđ šţéþš-çöḿþîļéđ ƒļöŵ çàŕŕîéš ţööļ ñöđéš öñļý, ƃéçàüšé ţĥé éñđš àŕé ƃîñđîñĝš (É-04) - Ñàḿé: ţĥé ŕéĝîšţéŕéđ ñàḿé öƒ ţĥé ţööļ (é.ĝ.
"pseudo-translate") - Ļàƃéļ: öþţîöñàļ đîšþļàý ļàƃéļ ƒöŕ ÜÎ ŕéñđéŕîñĝ
- Çöñƒîĝ: öþţîöñàļ ķéý-ṽàļüé çöñƒîĝüŕàţîöñ ḿàþ
- Þöšîţîöñ: ẋ/ý çööŕđîñàţéš ƒöŕ ṽîšüàļ ļàýöüţ îñ ţĥé ƒļöŵ éđîţöŕ
Ƃîñđîñĝš (É-04). À ƒļöŵ'š šöüŕçé àñđ šîñķ àŕé ƃîñđîñĝš ŕéšöļṽéđ ƒŕöḿ îñṽöçàţîöñ çöñţéẋţ (ƒîļé, ţĥé þŕöĵéçţ šţöŕé, à
.kpz, îñţéŕçĥàñĝé îḿþöŕţ/éẋþöŕţ, öŕ ñöñé), šö ţĥé šàḿé ƒļöŵ ŕüñš öṽéŕ àñý öŕîĝîñ. Éṽéŕý ƃüîļţ-îñ ƒļöŵ ĝŕàþĥ çàŕŕîéš ţööļ ñöđéš öñļý. Ţĥé ĝŕàþĥ îš çöḿþöšîţîöñ; à šîñĝļé ţööļ îš îñṽöķéđ đîŕéçţļý, ñöţ ŵŕàþþéđ îñ à öñé-ţööļ ƒļöŵ.
Éàçĥ FlowEdge çöññéçţš à šöüŕçé ñöđé ţö à ţàŕĝéţ ñöđé. TopologicalOrder()
çöḿþüţéš ţĥé éẋéçüţîöñ öŕđéŕ üšîñĝ Ķàĥñ'š àļĝöŕîţĥḿ, ŕéţüŕñîñĝ àñ éŕŕöŕ îƒ à
çýçļé îš đéţéçţéđ, šö îñṽàļîđ ƒļöŵ ĝŕàþĥš ñéṽéŕ ŕéàçĥ ţĥé ŕüñţîḿé éẋéçüţöŕ.
Ţĥé ƃüîļţ-îñ ƒļöŵ çàţàļöĝ îš à þŕöđüçţ çöñçéŕñ àñđ ļîṽéš îñ ţĥé ĥöšţ ḿöđüļé
(host/flowdef.BuiltInFlows), ñöţ îñ ţĥé éñĝîñé. Îţ çöṽéŕš ţŕàñšļàţîöñ ŵîţĥ
ĝüàŕđŕàîļš, çĥéçķš, ḿéḿöŕý ŕéüšé, þšéüđö-ţŕàñšļàţîöñ, ţĥé ḿéđîà-ţö-šüƃţîţļé
þàţĥš, àñđ ţĥé ŕéđàçţîöñ-ƃŕàçķéţéđ ƒļöŵš; ţĥé àüţĥöŕîţàţîṽé ļîšţ ŵîţĥ éàçĥ
ƒļöŵ'š šţéþš îš kapi flows àñđ ţĥé ĝéñéŕàţéđ
çöḿḿàñđ ŕéƒéŕéñçé.
kapi flows ļîšţš öñļý ţĥé çöḿþöšéđ (ḿüļţî-ţööļ) ƃüîļţ-îñ ƒļöŵš, ƃéçàüšé
šîñĝļé-ţööļ đéƒîñîţîöñš àŕé šüŕƒàçéđ àš ţöþ-ļéṽéļ ţööļ çöḿḿàñđš ŕàţĥéŕ ţĥàñ àš
ƒļöŵš. host/flowdef.FlowStore þéŕšîšţš üšéŕ-çŕéàţéđ ƒļöŵ đéƒîñîţîöñš àš ĴŠÖÑ
ƒîļéš öñ đîšķ, đîšţîñĝüîšĥéđ ƃý šöüŕçé:
built-in: šĥîþš ŵîţĥ neokapi, îḿḿüţàƃļéuser: çŕéàţéđ ƃý ţĥé üšéŕ, šţöŕéđ îñ ţĥé üšéŕ'š çöñƒîĝ đîŕéçţöŕý
À þŕöĵéçţ'š öŵñ ñàḿéđ ƒļöŵš àŕé đéçļàŕéđ îñ îţš ŕéçîþé'š flows: ƃļöçķ ŕàţĥéŕ
ţĥàñ šţöŕéđ àš ƒîļéš (É-04).
Šţéþš-ƃàšéđ ÝÀḾĻ ƒöŕḿàţ
À ĥüḿàñ-ƒŕîéñđļý šţéþš ƒöŕḿàţ îš ţĥé þŕîḿàŕý àüţĥöŕîñĝ šüŕƒàçé ƒöŕ ƒļöŵš îñ ÝÀḾĻ (É-03):
apiVersion: v1
kind: FlowDefinition
metadata:
name: Production Pipeline
spec:
steps:
- tool: recycle
config: { fuzzyThreshold: 75 }
- tool: translate
config: { provider: anthropic }
- tool: qa
Šţéþš àŕé šéǫüéñţîàļ ƃý đéƒàüļţ. parallel: ƃļöçķš þŕöṽîđé ƒàñ-öüţ. Ţĥé þàŕšéŕ
àüţö-đéţéçţš ţĥé šĥàþé (šţéþš ṽš ĝŕàþĥ) àñđ çöḿþîļéš šţéþš ţö ñöđéš àñđ éđĝéš.
Ƃöţĥ þŕöđüçé ţĥé šàḿé ŕüññàƃļé éẋéçüţöŕ.
Ţĥé šţéþš çàŕŕý öñļý ţĥé çöḿþöšîţîöñ. À ƒļöŵ'š šöüŕçé àñđ šîñķ àŕé ƃîñđîñĝš
ŕéšöļṽéđ àţ îñṽöçàţîöñ (ƒîļé, ţĥé þŕöĵéçţ šţöŕé, à .kpz, îñţéŕçĥàñĝé, öŕ ñöñé;
É-04) ŕàţĥéŕ ţĥàñ ƒîéļđš öƒ ţĥé ƒļöŵ đöçüḿéñţ.
Ƒàñ-öüţ àñđ ƃàţçĥîñĝ
tool.Tee() çöþîéš þàŕţš ţö Ñ öüţþüţ çĥàññéļš, éñàƃļîñĝ ƒàñ-öüţ ţöþöļöĝîéš ŵĥéŕé
öñé ñöđé ƒééđš ḿüļţîþļé þàŕàļļéļ ƃŕàñçĥéš. Ţĥé batch ţööļ çöļļéçţš ƃļöçķš îñţö
çöñƒîĝüŕàƃļé ƃàţçĥéš ƃéƒöŕé ƒöŕŵàŕđîñĝ, ŵĥîçĥ šüîţš ƃàţçĥ-çàþàƃļé ŕéḿöţé ÀÞΚ àñđ
ĻĻḾ þŕöḿþţš ţĥàţ ƃéñéƒîţ ƒŕöḿ ḿüļţîþļé îñþüţš þéŕ ŕéǫüéšţ
(Ḿ-05).
Šçŕîþţ šţéþ
Ţĥé script ţööļ ŕüñš üšéŕ-þŕöṽîđéđ ĴàṽàŠçŕîþţ (ÉŠ5) ṽîà ţĥé ĝöĵà ŕüñţîḿé. Éàçĥ
ţööļ îñšţàñçé öŵñš îţš öŵñ goja.Runtime, ŵĥîçĥ îš šàƒé ƃéçàüšé ToolFactory
ĝîṽéš öñé îñšţàñçé þéŕ ĝöŕöüţîñé. Ţĥé ĴŠ ÀÞÎ éẋþöšéš part, emit(), skip(),
àñđ log() ƒöŕ ƒîļţéŕîñĝ àñđ ţŕàñšƒöŕḿîñĝ þàŕţš: ļîĝĥţŵéîĝĥţ çüšţöḿ
ţŕàñšƒöŕḿàţîöñš ŵîţĥöüţ Ĝö çöđé. script îš îñ ţĥé éẋéç çļàšš: à ŕéçîþé
çàññöţ àŕḿ îţ šîļéñţļý (É-06).
Ţéŕḿîñöļöĝý: Okapi → neokapi
Ƒöŕ ŕéàđéŕš ƒàḿîļîàŕ ŵîţĥ ţĥé Okapi Ƒŕàḿéŵöŕķ, ţĥé éñĝîñé ḿàþš ţö Okapi çöñçéþţš àš ƒöļļöŵš:
| Okapi (Ĵàṽà) | neokapi (Ĝö) |
|---|---|
| Ƒîļţéŕ | ĐàţàƑöŕḿàţ (Ŕéàđéŕ/Ŵŕîţéŕ) |
| Šţéþ | Ţööļ |
| Þîþéļîñé | Ƒļöŵ |
| ÞîþéļîñéĐŕîṽéŕ | Éẋéçüţöŕ |
| Éṽéñţ | Þàŕţ |
| ŢéẋţÜñîţ | Ƃļöçķ |
| ŢéẋţƑŕàĝḿéñţ | Ŕüñ šéǫüéñçé ([]Run) |
| Çöđé | Ŕüñ |
| ŠţàŕţŠüƃĐöçüḿéñţ/ŠţàŕţŠüƃƑîļţéŕ | Çĥîļđ Ļàýéŕ |
Çöñšéǫüéñçéš
- Éàçĥ ţööļ ŕüñš çöñçüŕŕéñţļý; ḿüļţî-çöŕé ÇÞÜš àŕé üţîļîžéđ ŵîţĥîñ à šîñĝļé đöçüḿéñţ'š þîþéļîñé.
- Ḿüļţîþļé đöçüḿéñţš þŕöçéšš îñ þàŕàļļéļ, ƃöüñđéđ ƃý
MaxConcurrency. - Ƃàçķþŕéššüŕé îš àüţöḿàţîç: à šļöŵ ţööļ çàüšéš îţš îñþüţ çĥàññéļ ţö ƒîļļ, ŵĥîçĥ ƃļöçķš ţĥé üþšţŕéàḿ ţööļ ŵîţĥöüţ ḿàñüàļ çööŕđîñàţîöñ.
- Çöñţéẋţ çàñçéļļàţîöñ çļéàñļý þŕöþàĝàţéš ţĥŕöüĝĥ ţĥé éñţîŕé çĥàîñ.
- ŢööļƑàçţöŕîéš éñšüŕé ñö šĥàŕéđ ḿüţàƃļé šţàţé ƃéţŵééñ þàŕàļļéļ đöçüḿéñţš.
- Çöļļéçţöŕš þŕöṽîđé çŕöšš-đöçüḿéñţ àĝĝŕéĝàţîöñ ŵîţĥöüţ ƃŕéàķîñĝ ţĥé šţŕéàḿîñĝ ḿöđéļ.
- Ţööļ àüţĥöŕš đö ñöţ ḿàñàĝé ĝöŕöüţîñéš; ţĥé éẋéçüţöŕ ĥàñđļéš ļîƒéçýçļé, àñđ
ParallelBlockToolšüþþļîéš îñţŕà-ţööļ þàŕàļļéļîšḿ ŵîţĥöüţ àñý çöñçüŕŕéñçý çöđé îñ ţĥé ţööļ. StreamingCollectoréñàƃļéš ŕéàļ-ţîḿé öƃšéŕṽàţîöñ öƒ þîþéļîñé öüţþüţ ŵîţĥöüţ ḿöđîƒýîñĝ ţĥé Þàŕţ šţŕéàḿ öŕ àđđîñĝ ƃüƒƒéŕîñĝ šţàĝéš.- Ƒļöŵ ţŕàçîñĝ éñàƃļéš þöšţ-ĥöç đéƃüĝĝîñĝ àñđ ṽîšüàļîžàţîöñ, ĥéļþîñĝ üšéŕš üñđéŕšţàñđ ţööļ ƃéĥàṽîöüŕ àñđ îđéñţîƒý ƃöţţļéñéçķš.
TopologicalOrderṽàļîđàţîöñ çàţçĥéš çýçļéš ƃéƒöŕé ŕüñţîḿé, ĝîṽîñĝ ƒàšţ ƒééđƃàçķ đüŕîñĝ ƒļöŵ àüţĥöŕîñĝ.- ĴŠÖÑ àñđ ÝÀḾĻ šéŕîàļîžàţîöñ šüþþöŕţš îḿþöŕţ/éẋþöŕţ àñđ ṽéŕšîöñ çöñţŕöļ öƒ ƒļöŵ çöñƒîĝüŕàţîöñš.
- Šţéþš-ƃàšéđ ÝÀḾĻ ḿàķéš ƒļöŵ àüţĥöŕîñĝ àççéššîƃļé ŵîţĥöüţ Ĝö; ţĥé ṽîšüàļ éđîţöŕ àñđ ţĥé ÝÀḾĻ šţàý îñ šýñç ƃéçàüšé ƃöţĥ çöḿþîļé ţö ţĥé šàḿé ĝŕàþĥ.
Ŕéļàţéđ
- Ƒ-02: Ţĥé çöñţéñţ ḿöđéļ: ţĥé Þàŕţ ţýþéš ţĥàţ šţŕéàḿ
- É-02: Ţĥé ƒöŕḿàţ šýšţéḿ: ŕéàđéŕš ţĥàţ éḿîţ Þàŕţš, ŵŕîţéŕš ţĥàţ çöñšüḿé ţĥéḿ
- É-03: Ţĥé ţööļ šýšţéḿ: ţĥé ţööļš ţĥàţ ḿàķé üþ à ƒļöŵ
- É-04: Ƒļöŵš àñđ Î/Ö ƃîñđîñĝ: ŕéàđéŕ/ŵŕîţéŕ ƃéçöḿé šöüŕçé/šîñķ ƃîñđîñĝš; à ƒļöŵ îš çöḿþöšîţîöñ öñļý
- É-05: Ţĥé þļüĝîñ šýšţéḿ: þļüĝîñ ţööļš üšé ţĥé šàḿé éẋéçüţöŕ çöñţŕàçţ