I’m still confused a bit by reducing functions. Ultimately I am trying to do a (sort-by) on a core.async stream and am really strugging to get something to work. I feel like I am missing something fundamental. I’ll list a few things I have tried in the thread here.
One of the approaches I have tried is to use various permutations of net.cgrand.xforms, but its acting in ways that lead me to believe I dont really know how its supposed to be used.
For instance, this works as Id expect:
(eduction (net.cgrand.xforms/sort-by :foo) [{:foo "2"} {:foo "1"}])
=> ({:foo "1"} {:foo "2"})but when I try to use that xform on a core.async channel, I end up with essentially an infinitely repeating stream of sorted values, instead of terminating when the channel closes
I’ve tried various things such as using it as the xform parameter in clojure.core.async/chan, transduce, etc.
My real question is “how do I get a clojure.core/sort-by like behavior against a core.async stream. But I would also take explanations for what I am doing wrong with the cgrand library too, as it looks like there is some helpful transforms there I can use in other places if I learn how it is supposed to be used.
Other approaches I have taken is to forgo the cgrand library and write my own using clojure.core.async/transduce and a (sorted-map), but that didnt go great either
what are you trying to do? It doesn't really make sense to stream sorted values, because there's no way to know if you'll receive larger or smaller values in the future.
another approach would be collect the values using core.async/into, sorting the result and then putting those values onto another channel.
note that the input channel must be closed before a result is produced.
You'd expect the transducer to collect the values until the channel is closed and then return them sorted though no?
yea, there seems to be a bug in xforms sort-by.
in partition-all the completion arity first clears the state (https://github.com/clojure/clojure/blob/ee28f7f06d469b8e4f4eb48521ba51692146e774/src/clj/clojure/core.clj#L7403)
however, xform's sort (which is used by sort-by), does not (https://github.com/cgrand/xforms/blob/78076f8cd078ebb336ed9047df873ebc2ecf3aa1/src/net/cgrand/xforms.cljc#L557)
sorting as a step function is just nonsense anyway
Understood on the buffering limitation. What I am hoping to achieve is some degree of amortization where the sorting is being done concurrently to the upstream event generation, as the upstream has some expensive transformation work to do anyway.
Although, it seems like that should be unnecessary given > A completing process must call the completion operation on the final accumulated value exactly once. > https://clojure.org/reference/transducers#_creating_transducible_processes
Sorting is inherently a global process, it needs to see everything all at once
so maybe it's a bug in how core.async uses transducers since it's calling the completion arity more than once.
step functions are a streaming pass over the data
What I was experimenting with is having some of the overhead amortized while the stream is being generated so there is smaller latency once the stream closes
ok, seems to be a bug in core.async.
Yeah, that’s what I am doing now. It seems to work fine. I was more perplexed about why the cgrand/sort-by approach didn’t work
I am basically doing reduce and a mapcat xform after. I was originally trying to do it with cgrand or transduce, but hit weird problems that challenged my understanding
xforms is a mix of very useful parts, and experiments that are a bad idea
If you look at the implementation of the sort xform, there is no streaming amortization of anything
It just reduces into a mutable collection then sorts it at the end
The transducer could sort it at it's coming in though
Only if it still buffers everything up and doesn't emit anything to later steps in the pipeline until the end
Sorting is not a midstream operation, it is an operation at a stream sink, putting it in the middle of a pipeline stops the pipeline in its tracks
Understood. What I was after was incremental sorting, not some magical way to do streaming sort. In any case, I solved it with your suggestion of using async/reduce + sorted-map. Works well.
Thanks for the insights all!