【发布时间】:2017-04-27 15:12:02
【问题描述】:
我要做的是使用流从 mongodb 中提取一些数据,对数据进行一些操作,使用传入数据的键在 mongo 中进行另一个查询,然后加入数据。
实际代码要复杂得多,但我会用更简单的例子来说明。
我遇到的问题是,完成时永远不会被解雇,我假设是因为外部可观察对象永远不会完成。
const o$ = RxNode.fromStream(getDataFromMongo(), 'end'))
.filter(doSomeFiltering)
.bufferCount(100)
.flatMap(bufferedData => {
let ids = _.map(bufferedData, 'id');
return RxNode.fromStream(getMoreDataFromMongodb(ids))
.reduce(createMapIdObject(), {})
.flatMap(createdMap => { //because i get here, outer never complets
return Observable.from(bufferedData)
.map(item => ({
id: item.id,
size: createdMap[item.id].size
}));
});
})
.subscribe(
c => console.log(c),
err => console.log(err),
() => console.log('COMPLETE') // never happens
);
我怎样才能获得完整的射击?还是有更好的方法来完成我想要实现的目标?
如果我在嵌套的flatMap 内将.finally(() => console.log(done) 链接到Observable.from,则会被解雇。所以内部 observable 正在完成,但外部没有。
【问题讨论】:
-
一些 observables 不会发出完整的通知。为什么在成功函数中不能做自己想做的事情
-
不确定我是否理解您在成功功能中执行此操作的意思。当我达到成功功能时,我只有一组数据,我需要将它们结合起来。
-
你能把整个东西包装在一个 observable 中吗?
-
getDataFromMongo 是什么样子的?
-
这是一个流:
db.coll.find({some: conditions}, {some: fields}).stream();
标签: javascript mongodb rxjs rxjs5