如何用RxJS正确发起嵌套HATEOAS请求?附StackBlitz示例
解决方案:RxJS关联OGC SensorThings API嵌套资源并定时更新动态数据
一、关联Thing与Datastream资源
由于OGC SensorThings API采用HATEOAS规范,资源通过导航URI关联而非外键,可通过RxJS的mergeMap+forkJoin组合实现资源关联:
- 先获取所有Thing列表
- 对每个Thing,通过其
Datastreams@iot.navigationLink发起请求获取关联的Datastream集合 - 将Datastream合并到对应的Thing对象中,生成完整关联对象
代码示例:
import { forkJoin, map, mergeMap } from 'rxjs'; // 封装API请求方法(根据实际项目调整) const fetchThings = () => fetch('/api/Things').then(res => res.json()); const fetchResourceByUri = (uri: string) => fetch(uri).then(res => res.json()); // 获取带关联Datastream的完整Thing集合 const getThingsWithDatastreams = () => { return fetchThings().pipe( mergeMap(things => { // 并行处理每个Thing的Datastream请求 const thingRequests$ = things.map(thing => { return fetchResourceByUri(thing['Datastreams@iot.navigationLink']).pipe( map(datastreams => ({ ...thing, datastreams })) ); }); return forkJoin(thingRequests$); }) ); };
关键说明:
mergeMap用于将Thing列表转换为多个Datastream请求的Observable流forkJoin确保所有Datastream请求完成后,一次性返回合并后的完整Thing集合- 注意API返回的导航属性格式(如
Datastreams@iot.navigationLink),需与实际响应字段匹配
二、定时更新动态资源(Observations、Locations)
对于Observations这类动态数据,可通过RxJS的interval+BehaviorSubject实现定时刷新:
- 用
BehaviorSubject保存当前的完整Thing集合状态,便于组件订阅和更新 - 初始化时加载关联好的完整数据
- 每隔固定时间重新请求动态资源,更新状态中的对应字段
代码示例:
import { interval, switchMap, BehaviorSubject } from 'rxjs'; // 初始化状态Subject,保存完整Thing集合 const thingsState$ = new BehaviorSubject<any[]>([]); // 首次加载关联数据 getThingsWithDatastreams().subscribe(initialThings => { thingsState$.next(initialThings); }); // 定时更新(示例为每30秒刷新一次Observations) const REFRESH_INTERVAL = 30000; interval(REFRESH_INTERVAL).pipe( switchMap(() => { const currentThings = thingsState$.value; // 对每个Thing的Datastream重新请求Observations const updatedThingRequests$ = currentThings.map(thing => { const updatedDatastreams$ = thing.datastreams.map(datastream => { return fetchResourceByUri(datastream['Observations@iot.navigationLink']).pipe( map(observations => ({ ...datastream, observations })) ); }); return forkJoin(updatedDatastreams$).pipe( map(updatedDatastreams => ({ ...thing, datastreams: updatedDatastreams })) ); }); return forkJoin(updatedThingRequests$); }) ).subscribe(updatedThings => { thingsState$.next(updatedThings); }); // 组件中订阅状态获取实时数据 thingsState$.subscribe(things => { // 在这里处理组件渲染逻辑 });
关键说明:
BehaviorSubject用于维护全局状态,确保组件能获取最新数据switchMap用于取消上一次未完成的更新请求,避免竞态问题- 可根据需求调整刷新间隔,或针对不同动态资源(如Locations)单独实现更新逻辑
内容的提问来源于stack exchange,提问作者Idea
相关产品推荐
相关产品推荐

