Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- rddQueue = ssc.queueStream([rdd1, rdd2])
- def func(new_values, old_value):
- return sum(new_values) + (old_value or 0)
- rddQueue = rddQueue.updateStateByKey(func).transform(lambda x: x.sortBy(lambda y: y[1], ascending=False))
- .updateStateByKey(func).transform(...
Add Comment
Please, Sign In to add comment