【发布时间】:2018-10-16 04:42:59
【问题描述】:
我正在尝试从一组项目中创建一个 Observable,每个项目都会定期检查服务器更新,然后在获得每个项目所需的结果时发送一个操作。
下面的答案很有帮助,但不是我想要的
这是我一直在尝试的另一种方法:
export function handleProcessingScenes(action$,store) {
return action$.ofType(REQUEST_ALL_SCENES_BY_LOCATION_FULFILLED)
.switchMap(({ scenesByLocation }) => Observable.from(scenesByLocation))
.filter(scene => scene.scenePanoTask)
.mergeMap(scene => updateScene(scene))
}
function updateScene(scene) {
return Observable.interval(3000)
.flatMap(() => requestSceneUpdates(scene.id))
.takeWhile(res => res.payload.status < 4)
.timeout(600000, Observable.throw(new Error('Timeout')))
}
API 函数返回一个 Observable
export function requestSceneUpdates(sceneId){
console.log('requestSceneUpdate')
const request = fetch(`${API_URL}/scene/task/${sceneId}/update`, {
method: 'get',
credentials: 'include',
crossDomain: true,
}).then(res => res.json())
return Observable.fromPromise(request)
}
但是,这只会调用一次“requestSceneUpdate”函数。
我基本上想为场景中的每个场景每 3 秒调用一次该函数。然后我想在每个动作完成后返回一个动作。
我对一个场景的史诗是
export function sceneProcessingUpdate(action$) {
return action$.ofType(REQUEST_SCENE_PROCESSING_TASK_SUCCESS)
.switchMap(({task}) =>
Observable.timer(0, 30000).takeUntil(action$.ofType( REQUEST_SCENE_PROCESSING_TASK_UPDATE_SUCCESS))
.exhaustMap(() =>
requestSceneUpdates(task.id)
.map((res) => {
if (res.error)
return { type: REQUEST_SCENE_PROCESSING_TASK_UPDATE_FAILED, message: res.message }
else if(res.payload.status === 4)
return { type: REQUEST_SCENE_PROCESSING_TASK_UPDATE_SUCCESS, task: res.payload }
else
return requestSceneProcessingTaskMessage(res.payload)
})
.catch(err => { return { type: REQUEST_SCENE_PROCESSING_TASK_UPDATE_FAILED, message: err } })
)
)
}
【问题讨论】:
-
我认为描述您想要实现的目标和输入内容会更容易。
-
谢谢,我添加了更多细节。
-
所以您希望
requestSceneUpdates在成功后每 3 秒停止一次发射?但只针对那个场景还是针对所有场景?在我看来,我会使用重试机制,而不是计时器。因此,如果失败,请在 3 秒后重试(直到不再失败)... -
看起来问题是您没有在
catch的回调中返回 Observable。你应该用例如包装它。Observable.of()。否则发布整个错误消息。
标签: rxjs observable