feat(wirings): Kafka Streams DSL tier wirings + capstone demo (stacks on #3; Kafka mirror of Flink #6→#10) - #4
Closed
estebanzimanyi wants to merge 2 commits into
Conversation
… surface (Kafka mirror of MobilityFlink MobilityDB#5) Mirror of MobilityFlink codegen/flink-meos-ops (PR MobilityDB#5) for the Kafka Streams binding. Same generators, same tier classification, same catalog source — differs only in package path (`org.mobilitydb.kafka.meos`) and module layout (`kafka-streams-app/` vs `flink-processor/`). Adds 50 `MeosOps<Class>` + 6 `MeosOpsFree<Header>` + 1 shared `MeosOpsRuntime` = 57 Java classes, 2,097 methods (77.7% of JMEOS PR MobilityDB#19's 2,699-method surface). Stacks on feat/jmeos-bridge-swap; additive-only; touches no existing file. See MobilityFlink MobilityDB#5 for the full tier vocabulary, regeneration recipe, and coexistence design with `berlinmod.MEOSBridge`.
…facades + capstone demo Adds the org.mobilitydb.kafka.meos.wirings package — Kafka mirror of the MobilityFlink wirings (PRs MobilityDB#6→MobilityDB#10). Single PR with all four tiers + runnable composite demo, since Kafka Streams' DSL is naturally lambda-driven and most tier wirings collapse to small static-factory classes returning serializable functional-interface implementations. ## Files - MeosStatelessOps.java — predicate/intPredicate/mapper factories returning Predicate<K,V> / ValueMapper<V,R> for KStream.filter / .mapValues (covers stateless tier: 804 methods + io-meta: 195 methods) - MeosBoundedStateProcessor.java — full Processor<KIn,VIn,KOut,VOut> class with KeyValueStore<KIn,byte[]> for per-key MEOS-handle state that survives changelog replay / rebalance (covers bounded-state tier: 797 methods) - MeosWindowedAggregator.java — initializer/aggregator factories for KStream.groupByKey().windowedBy(...).aggregate(...) (covers windowed tier: 161 methods) - MeosCrossStreamJoiner.java — joiner factory wrapping a serializable ValueJoiner for KStream.join(other, joiner, JoinWindows) (covers cross-stream tier: 140 methods) - MeosOpsRuntime.java — wirings-package alias for the codegen package's MeosOpsRuntime.MEOS_AVAILABLE flag (libmeos probed once per JVM) - demo/MeosWiringsDemoTopology.java — runnable Kafka Streams topology composing all four tier wirings into one pipeline (per-region running-union → 30s tumbling aggregate → ±1m cross-stream join against region-queries); main() prints topology description always, instantiates TopologyTestDriver when MEOS_AVAILABLE (no broker required) - README.md — tier vocabulary, lambda-first design rationale, full recipe + demo walkthrough, coexistence with berlinmod.MEOSBridge ## Cumulative wirings-layer coverage Same as the Flink side: 2,097 of 2,097 emitted methods (100%) wirable through 5 generic classes; zero per-method registration. ## Design choice — lambda-first Kafka Streams' DSL accepts lambdas directly (Predicate, ValueMapper, Aggregator, ValueJoiner). Only bounded-state needs a real class for state-store binding via Processor.init(ProcessorContext). The wirings reflect that asymmetry: small static-factory classes for the four lambda-shaped tiers, one full Processor class for bounded-state. Adopters wanting a class-shaped wiring (matching Flink for cross-binding parity) can subclass any of the helpers — the serializable functional interfaces are public. ## Stacks on PR MobilityDB#3 (codegen mirror) Additive-only: 6 new files under kafka-streams-app/src/main/java/org/mobilitydb/kafka/meos/wirings/. Touches no existing file. Locally compile-verified: 110 .class files total (94 from PR MobilityDB#3 base + 16 new — 4 wiring classes + 11 nested lambda interfaces + MeosOpsRuntime + 1 demo class).
Member
Author
|
Consolidated into #13 (consolidate/kafka-streaming-implementation). |
Member
Author
|
Superseded by the Path-B consolidation: the former 12-deep stack is collapsed into three reviewable topical PRs — scaffold #1 → MEOS integration #14 → benchmark #15 — each one clean squashed commit with the generated-facade bulk, dead family-flag profiles, and invented synthetic corpus removed. Closing as folded into #14/#15. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Kafka-side mirror of the MobilityFlink wirings stack (PRs #6→#10). Single PR with all four tier wirings + runnable composite demo, since Kafka Streams' DSL is naturally lambda-driven and most tier wirings collapse to small static-factory classes returning serializable functional-interface implementations.
Stacks on PR #3 (the Kafka codegen mirror). Additive-only — no existing file touched.
What's in this PR (6 new files under
kafka-streams-app/src/main/java/org/mobilitydb/kafka/meos/wirings/)stateless(804 methods)MeosStatelessOpspredicate(...)/intPredicate(...)/mapper(...)→Predicate<K,V>/ValueMapper<V,R>forKStream.filter/.mapValuesbounded-state(797 methods)MeosBoundedStateProcessorProcessor<KIn,VIn,KOut,VOut>class withKeyValueStore<KIn,byte[]>for per-key MEOS-handle state that survives changelog replay / rebalancewindowed(161 methods)MeosWindowedAggregatorinitializer(...)/aggregator(...)forKStream.groupByKey().windowedBy(...).aggregate(...)cross-stream(140 methods)MeosCrossStreamJoinerjoiner(...)wrapping serializableValueJoinerforKStream.join(other, joiner, JoinWindows), same-key pairing, time-bounded match windowio-meta(195 methods)MeosStatelessOps.mapper(...)sequence-only(14 methods)Plus:
MeosOpsRuntime.java— wirings-package alias for the codegen package'sMEOS_AVAILABLEflagdemo/MeosWiringsDemoTopology.java— runnable Kafka Streams topology composing all 4 tier wirings; usesTopologyTestDriver(no broker required)README.md— full tier vocabulary, lambda-first design rationale, demo walkthrough, coexistence withberlinmod.MEOSBridgeCumulative wirings coverage: 2,097 of 2,097 emitted methods (100%) wirable through 5 generic classes; zero per-method registration. Identical to Flink side.
Design choice — lambda-first
Kafka Streams' DSL accepts lambdas directly (
Predicate,ValueMapper,Aggregator,ValueJoiner). Onlybounded-stateneeds a real class for state-store binding viaProcessor.init(ProcessorContext). The wirings reflect that asymmetry — small static-factory classes for the four lambda-shaped tiers, one fullProcessorclass forbounded-state.Adopters wanting a class-shaped wiring (matching Flink for cross-binding parity) can subclass any of the helpers — the serializable functional interfaces are public.
State-store discipline (bounded-state)
Same as the Flink wirings: raw
jnr.ffi.Pointerdoes not survive Kafka Streams' changelog-replay / rebalance / state-rebuild paths (the state-store changelog is a Kafka topic; state must be byte-serializable to be replayable).MeosBoundedStateProcessorstores state asbyte[](typically MEOS-WKB or MEOS-WKT) with three adopter-supplied lambdas mediating the round-trip:First record for a key sees
prior == null(no prior state); wiring skips deserialize and lets step seed.End-to-end demo
MeosWiringsDemoTopologycomposes all four tier wirings into one Kafka Streams topology:vehicle-eventssource → 2. stateless filter → 3. bounded-state processor (per-region running tbox union) → 4. windowed aggregator (30s tumbling) → 5. cross-stream joiner (±1m bound, joined againstregion-queriestopic) → 6.overlap-outputsink.Run via
mvn exec:java.main()always printsTopology.describe(); when MEOS available, instantiatesTopologyTestDriver(no broker required).Compile verification
Locally green: 110 .class files total (94 from PR #3 base + 16 new — 4 wiring classes + 11 nested lambda interfaces +
MeosOpsRuntime+ 1 demo class).Stacking
Stacks on
codegen/kafka-meos-ops(PR #3). Whole 11-file diff is contained underkafka-streams-app/src/main/java/org/mobilitydb/kafka/meos/wirings/.