Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

[#238] Add divert func #236

Open
wants to merge 7 commits into
base: master
Choose a base branch
from
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
## [0.7.8] - [2021-03-01]
* Override the netty version pulled by Aleph with one which fixes https://nvd.nist.gov/vuln/detail/CVE-2020-11612 [#261](https://github.com/FundingCircle/jackdaw/pull/261)
* Restore the test fixture namespace [#266](https://github.com/FundingCircle/jackdaw/pull/266)
* Add `divert?` stream operation

## [0.7.7] - [2021-02-09]

Expand Down
14 changes: 14 additions & 0 deletions src/jackdaw/streams.clj
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,20 @@
[kstream predicate-fns]
(p/branch kstream predicate-fns))

(defn divert?
"Diverts records that match any `pred` to the provided `topic-config`. Records that
do not match any `pred` are pushed on through the app.

When providing multiple diverts, diverts can either be:
- ordered (by providing a vector of `predicate` `topic-config` tuples); or
- unordered (by providing a map `predicate` to `topic-config` key pairs)"
([stream diverts]
(clojure.core/reduce #(apply divert? %1 %2) stream (into [] diverts)))
([stream pred topic-config]
(let [[divert-stream continue-stream] (branch stream [pred (constantly true)])]
(to divert-stream topic-config)
continue-stream)))

(defn flat-map
"Creates a KStream that will consist of the concatenation of messages
returned by calling `key-value-mapper-fn` on each key/value pair in the
Expand Down
72 changes: 72 additions & 0 deletions test/jackdaw/streams_test.clj
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,78 @@
(is (= [[1 1]] (mock/get-keyvals driver topic-pos)))
(is (= [[1 -1]] (mock/get-keyvals driver topic-neg)))))

(testing "divert?"
(let [topic-a (mock/topic "topic-a")
topic-odd (mock/topic "topic-odd")
topic-even (mock/topic "topic-even")
driver (mock/build-driver (fn [builder]
(-> builder
(k/kstream topic-a)
(k/divert? (comp odd? last) topic-odd)
(k/to topic-even))))
publish (partial mock/publish driver topic-a)]

(publish 1 1)
(publish 1 2)
(publish 1 3)
(publish 1 4)

(is (= [1 3] (map last (mock/get-keyvals driver topic-odd))))
(is (= [2 4] (map last (mock/get-keyvals driver topic-even)))))
(testing "multiple diverts"
(testing "un-ordered sorting"
(let [topic-a (mock/topic "topic-a")
topic-four (mock/topic "topic-mod-4")
topic-five (mock/topic "topic-mod-5")
topic-seven (mock/topic "topic-mod-7")
topic-rest (mock/topic "topic-rest")

mod-four? (comp zero? #(mod % 4) last)
mod-five? (comp zero? #(mod % 5) last)
mod-seven? (comp zero? #(mod % 7) last)

driver (mock/build-driver (fn [builder]
(-> builder
(k/kstream topic-a)
(k/divert? {mod-four? topic-four
mod-five? topic-five
mod-seven? topic-seven})
(k/to topic-rest))))
publish (partial mock/publish driver topic-a)]

(doseq [i (range 10)]
(publish 1 (inc i)))

(is (= [4 8] (map last (mock/get-keyvals driver topic-four))))
(is (= [5 10] (map last (mock/get-keyvals driver topic-five))))
(is (= [7] (map last (mock/get-keyvals driver topic-seven))))
(is (= [1 2 3 6 9] (map last (mock/get-keyvals driver topic-rest))))))

(testing "ordered sorting"
(let [topic-a (mock/topic "topic-a")
topic-three (mock/topic "topic-mod-3")
topic-five (mock/topic "topic-mod-5")
topic-rest (mock/topic "topic-rest")

mod-three? (comp zero? #(mod % 3) last)
mod-five? (comp zero? #(mod % 5) last)

driver (mock/build-driver (fn [builder]
(-> builder
(k/kstream topic-a)
(k/divert? [[mod-five? topic-five]
[mod-three? topic-three]])
(k/to topic-rest))))
publish (partial mock/publish driver topic-a)]

(doseq [i (range 15)]
(publish 1 (inc i)))

(is (= [3 6 9 12] (map last (mock/get-keyvals driver topic-three))))
;; mod-5 gets 15 b/c it was applied first
(is (= [5 10 15] (map last (mock/get-keyvals driver topic-five))))
(is (= [1 2 4 7 8 11 13 14] (map last (mock/get-keyvals driver topic-rest))))))))

(testing "flat-map"
(let [topic-a (mock/topic "topic-a")
topic-b (mock/topic "topic-b")
Expand Down