Skip to content

Instantly share code, notes, and snippets.

@chitacan
Last active March 11, 2016 12:00
Show Gist options
  • Select an option

  • Save chitacan/88577b8f1f8dc5df4d18 to your computer and use it in GitHub Desktop.

Select an option

Save chitacan/88577b8f1f8dc5df4d18 to your computer and use it in GitHub Desktop.
console.clear()
# Recipes from "RXJS Lessons"
# (https://www.youtube.com/playlist?list=PLX7ZnQs0Gb1pc1vNduNaIgFiyQ3THgymv)
# (https://jsbin.com/garaci/31/edit)
# ----------------------------------------------
# 1. creating an observable
src_1 = Rx.Observable.create (observer) ->
id = setTimeout () ->
console.log 'timeout hit'
observer.onNext 42
observer.onCompleted()
,1000
console.log 'started'
() ->
clearTimeout id
console.log 'disposed'
sub_1 = src_1.subscribe (x) ->
console.log "next #{x}"
, (err) ->
console.error err
, () -> console.log 'done'
setTimeout () ->
sub_1.dispose()
, 500
# ----------------------------------------------
# 2. what is rxjs
src_2 = [0,1,2,3,4,5]
src_2 = Rx.Observable.fromArray [0,1,2,3,4,5]
src_2 = Rx.Observable.interval 1000
.take 6
src_2.filter (x) -> x % 2 is 1
.map (x) -> x + '!'
.forEach (x) -> console.log x
# ----------------------------------------------
# 3. observables vs promises
promise = new Promise (resolve) ->
setTimeout () ->
resolve 42
, 1000
# this will hit even if comment out 'promise.then'...
console.log 'promise started'
# promise.then (x) -> console.log "promise #{x}"
src_3 = Rx.Observable.create (observer) ->
setTimeout () ->
observer.onNext 42
, 1000
console.log 'observable started'
# src_3.subscribe (x) -> console.log "observable #{x}"
# ----------------------------------------------
# 4. throttled buffering in rxjs
btn = document.querySelector '#click'
click = Rx.Observable.fromEvent btn, 'click'
open = Rx.Observable.interval 1000
sendValues = (arr) ->
pre = document.createElement 'pre'
pre.innerHTML = JSON.stringify arr
document.querySelector '#result'
.appendChild pre
click.scan (x) ->
x + 1
, 0
#.buffer open
.buffer click.throttle 1000
.filter (x) -> x.length > 0
.forEach (x) -> sendValues x
# ----------------------------------------------
# 5. stream processing with rxjs vs array higher order functions
src_5 = [0,1,2,3,4,5]
result = src_5.filter (x, i, arr) ->
console.log "filtering : #{x}"
console.log "filter : #{arr is src_5}"
x % 2 is 0
.map (x, i, arr) ->
console.log "mapping : #{x}"
console.log "map : #{arr is src_5}"
x + '!'
, ''
.reduce (r, x, i, arr) ->
console.log "reducing : #{x}"
console.log "reduce : #{arr is src_5}"
r + x
console.log result
src_5 = Rx.Observable.fromArray [0,1,2,3,4,5]
src_5.filter (x) ->
console.log "filtering : #{x}"
x % 2 is 0
.map (x) ->
console.log "mapping : #{x}"
x + '!'
.reduce (r, x) ->
console.log "reducing : #{x}"
r + x
, ''
.subscribe (x) -> console.log x
# ----------------------------------------------
# 6. toggle a stream on / off with rxjs
result = document.querySelector '#result'
toggle = document.querySelector '#toggle'
checked = Rx.Observable.fromEvent toggle, 'change'
.map (e) -> e.target.checked
src_6 = Rx.Observable.interval 500
.map (x) -> '.'
checked.filter (x) -> x is yes
.flatMapLatest -> src_6.takeUntil checked
.subscribe (x) ->
result.innerText += x
# ----------------------------------------------
# 7. map vs flatmap
src_7 = Rx.Observable.interval 1000
.take 10
# .map (x) -> Rx.Observable.timer(5000).map () -> x * 2
# .mergeAll()
.flatMap (x) -> Rx.Observable.timer(5000).map () -> x * 2
src_7.subscribe (x) -> console.log x
# ----------------------------------------------
# 8. demystifying cold and hot rxjs observables
clock = Rx.Observable.interval 1000
.take 10
.map (x) -> x + 1
.publish().refCount()
setTimeout () ->
clock.subscribe (x) -> console.log "a : #{x}"
, 1000
setTimeout () ->
clock.subscribe (x) -> console.log " b : #{x}"
, 4000
# ----------------------------------------------
# 9. aggregating streams with reduce and scan
# src_9 = Rx.Observable.fromArray [0,1,2,3,4,5]
src_9 = Rx.Observable.interval 1000
.take 5
# src_9.reduce (r, x) ->
src_9.scan (r, x) ->
r + x
, 0
.subscribe (x) -> console.log x
# ----------------------------------------------
# 10. introduction to connectable observable and using publish refcount
clock = Rx.Observable.interval 1000
.take 10
.map (x) -> x + 1
.startWith 0
.publish().refCount()
clock.subscribe (x) -> console.log "a : #{x}"
setTimeout () ->
clock.subscribe (x) -> console.log " b : #{x}"
, 3000
# ----------------------------------------------
# 11. error handling
Rx.Observable.of 1,2,3,4
.map (x) ->
if x is 3
throw 'i hate three'
else
x
# .onErrorResumeNext Rx.Observable.just('go ahead !')
# .catch (err) -> Rx.Observable.just('go ahead')
# .retry 3
.retryWhen (err) ->
err.delay 1000
.take 5
.concat Rx.Observable.throw 'fuck'
.subscribe (x) ->
console.log x
, (err) ->
console.error err
, () ->
console.log 'done'
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment