这篇翻译之后英文原文有更新,内容可能已过时。
Manager 与 Middleware
Reactive Data Client 采用 flux store 模式,其特点是 store 的单向数据流易于理解和调试。状态更新由 reducer 函数执行。

在 flux 架构中,flux 循环里的所有函数都必须是 纯函数,这一点至关重要。 Manager 负责集中编排副作用。换句话说,它们是 Data Client 与外部世界交互的途径。
例如,NetworkManager 负责编排数据获取,SubscriptionManager 则跟踪哪些资源通过 useLive 或 useSubscription 被订阅。通过集中控制,NetworkManager 会自动对 fetch 去重,而 SubscriptionManager 只会让正在渲染的资源保持更新。
因此,Manager 是集成其他副作用的最佳方式,例如日志、错误上报、指标、通知、数据流、聚焦或重连时刷新、跨标签页同步以及离线持久化。也可以通过自定义它们来改变核心行为。
| 默认 Manager | |
|---|---|
| NetworkManager | 将 fetch dispatch 转换为网络调用 |
| SubscriptionManager | 处理轮询订阅 |
| DevToolsManager | 支持调试 |
| 额外的 Manager | |
| LogoutManager | 处理 HTTP 401(或其他登出条件) |
示例
Reactive Data Client 通过其 Controller 执行 dispatch 和 store 访问,从而提升类型安全性和易用性
Middleware 日志
import type { Manager, Middleware } from '@data-client/react';
export default class LoggingManager implements Manager {
middleware: Middleware = controller => next => async action => {
console.log('before', action, controller.getState());
await next(action);
console.log('after', action, controller.getState());
};
cleanup() {}
}
错误上报
通过检查设置了 error 的 SET_RESPONSE action,将失败的 fetch 上报给 Sentry 等监控服务。
import {
type Manager,
type Middleware,
actionTypes,
} from '@data-client/react';
import { captureException } from '@sentry/react';
export default class ErrorReportManager implements Manager {
middleware: Middleware = controller => next => async action => {
if (action.type === actionTypes.SET_RESPONSE && action.error)
captureException(action.response, {
extra: { endpoint: action.endpoint.name, args: action.args },
});
return next(action);
};
cleanup() {}
}
指标
通过观察 FETCH action 来跟踪 fetch 耗时。action.meta.promise
会在 fetch 完成时 resolve。
import {
type Manager,
type Middleware,
actionTypes,
} from '@data-client/react';
import { trackTiming } from './analytics';
export default class MetricsManager implements Manager {
middleware: Middleware = controller => next => async action => {
if (action.type === actionTypes.FETCH) {
const start = performance.now();
action.meta.promise
.finally(() => {
trackTiming(action.endpoint.name, performance.now() - start);
})
// the fetch's caller handles errors; this only observes timing
.catch(() => {});
}
return next(action);
};
cleanup() {}
}
通知(toast)
在任意变更成功或失败时显示 toast。
import {
type Manager,
type Middleware,
actionTypes,
} from '@data-client/react';
import { toast } from './toast';
export default class ToastManager implements Manager {
middleware: Middleware = controller => next => async action => {
if (
action.type === actionTypes.SET_RESPONSE &&
action.endpoint.sideEffect
) {
if (action.error) toast.error(`${action.endpoint.name} failed`);
else toast.success(`${action.endpoint.name} succeeded`);
}
return next(action);
};
cleanup() {}
}
聚焦或重连时刷新
Controller.expireAll() 会将数据标记为过时,从而在不挂起的情况下触发所有_正在渲染_的数据重新获取(stale-while-revalidate)。 init() 和 cleanup() 负责管理事件监听器。
import type { Manager, Middleware, Controller } from '@data-client/react';
export default class RefreshManager implements Manager {
declare protected controller: Controller;
protected handle = () =>
this.controller.expireAll({ testKey: () => true });
middleware: Middleware = controller => {
this.controller = controller;
return next => async action => next(action);
};
init() {
window.addEventListener('focus', this.handle);
window.addEventListener('online', this.handle);
}
cleanup() {
window.removeEventListener('focus', this.handle);
window.removeEventListener('online', this.handle);
}
}
跨标签页同步
当某个标签页中的变更成功时,使用 BroadcastChannel 将其他所有标签页中的数据标记为过时。
import {
type Manager,
type Middleware,
actionTypes,
} from '@data-client/react';
export default class TabSyncManager implements Manager {
protected channel = new BroadcastChannel('data-client');
middleware: Middleware = controller => {
this.channel.onmessage = () =>
controller.expireAll({ testKey: () => true });
return next => async action => {
if (
action.type === actionTypes.SET_RESPONSE &&
action.endpoint.sideEffect &&
!action.error
)
this.channel.postMessage('mutation');
return next(action);
};
};
cleanup() {
this.channel.close();
}
}
离线持久化
使用 IndexedDB 持久化 store
(这里通过 idb-keyval);并通过
DataProvider 的 initialState 恢复它。IndexedDB 的写入是异步的,并且使用结构化克隆,而不会像 localStorage 那样因 JSON 序列化而阻塞主线程。对写入进行防抖可以让密集的 action 突发保持低开销。恢复时请考虑过期时间。
import type { Manager, Middleware } from '@data-client/react';
import { set } from 'idb-keyval';
export default class PersistManager implements Manager {
declare protected timer?: ReturnType<typeof setTimeout>;
middleware: Middleware = controller => next => async action => {
await next(action);
// debounce: persist at most once per second
clearTimeout(this.timer);
this.timer = setTimeout(() => {
// in-flight optimistic updates reference functions, so are not persistable
const state = { ...controller.getState(), optimistic: [] };
set('data-client', state);
}, 1000);
};
cleanup() {
clearTimeout(this.timer);
}
}
import { DataProvider, getDefaultManagers } from '@data-client/react';
import { createRoot } from 'react-dom/client';
import { get } from 'idb-keyval';
import App from './App';
import PersistManager from './PersistManager';
const managers = [...getDefaultManagers(), new PersistManager()];
const initialState = await get('data-client');
createRoot(document.body).render(
<DataProvider initialState={initialState} managers={managers}>
<App />
</DataProvider>,
);
Middleware 数据流(基于推送)
添加一个 Manager 来处理服务器通过 websockets 或 Server Sent Events 推送的数据,可以确保在数据更新与用户操作无关时,我们依然能保持数据新鲜。例如交易应用中的价格,或实时协作编辑器。
import type {
Manager,
Middleware,
Controller,
EntityInterface,
} from '@data-client/react';
export default class StreamManager implements Manager {
declare protected controller: Controller;
declare protected evtSource: WebSocket | EventSource;
declare protected createEventSource: () => WebSocket | EventSource;
declare protected entities: Record<string, EntityInterface>;
constructor(
createEventSource: () => WebSocket | EventSource,
entities: Record<string, EntityInterface>,
) {
this.createEventSource = createEventSource;
this.entities = entities;
}
middleware: Middleware = controller => {
this.controller = controller;
return next => async action => next(action);
};
connect() {
this.evtSource = this.createEventSource();
this.evtSource.onmessage = (event: MessageEvent) => {
try {
const msg: { type: string; args: [any]; data: any } = JSON.parse(
event.data,
);
if (msg.type in this.entities)
this.controller.set(
this.entities[msg.type],
...msg.args,
msg.data,
);
} catch (e) {
console.error('Failed to handle message');
console.error(e);
}
};
}
init() {
this.connect();
}
cleanup() {
this.evtSource?.close();
}
}
Controller.set() 允许直接使用 event.data
更新可查询 schema。
批量处理高频更新
像交易所行情这样的数据流每秒可能发送数百条消息,而且连接建立时通常会先发送一个大型快照。与其对每条消息都调用 set(),不如先将它们缓冲起来,再使用 Array schema 按批写入。
Controller.set([Entity], rows) 会在一次 store 更新中规范化所有行。
export default class StreamManager implements Manager {
// ...
protected buffer: Record<string, any[]> = {};
declare protected flushTimeout?: ReturnType<typeof setTimeout>;
connect() {
this.evtSource = this.createEventSource();
this.evtSource.onmessage = event => {
const msg = JSON.parse(event.data);
if (msg.type in this.entities) {
(this.buffer[msg.type] ??= []).push(msg.data);
this.flushTimeout ??= setTimeout(this.flush, 50);
}
};
}
flush = () => {
const buffer = this.buffer;
this.buffer = {};
this.flushTimeout = undefined;
for (const type in buffer) {
this.controller.set([this.entities[type]], buffer[type]);
}
};
cleanup() {
this.evtSource?.close();
clearTimeout(this.flushTimeout);
this.flushTimeout = undefined;
this.buffer = {};
}
}
同一批次中 pk 相同的行会按顺序合并,并跳过 Entity.shouldReorder(),因此当顺序很重要时,每个 pk 只缓冲最新的一条消息。
试试下面的两个按钮。这个浏览器测试从空的 store 开始,分别计时 500 次 set()
调用的 Promise.all 与一次批量 set()。两种方式都只产生一次 React 提交,并且各自写入 500 个新价格。
import React from 'react'; import { useController, useQuery } from '@data-client/react'; import { Ticker, newPrices } from './Ticker'; function PriceStream() { const ctrl = useController(); const [timing, setTiming] = React.useState(''); const first = useQuery(Ticker, { product_id: 'COIN-0' }); const time = async ( label: string, write: (rows: ReturnType<typeof newPrices>) => Promise<unknown>, ) => { const rows = newPrices(); const start = performance.now(); await write(rows); setTiming(`${label}: ${(performance.now() - start).toFixed(1)} ms`); }; const perRow = () => time('500 set() calls', rows => Promise.all( rows.map(row => ctrl.set(Ticker, { product_id: row.product_id }, row)), ), ); // highlight-next-line const batch = () => time('1 batch set()', rows => ctrl.set([Ticker], rows)); return ( <div> <button onClick={perRow}>set() per row</button>{' '} <button onClick={batch}>batch set()</button> <p>COIN-0: {first ? `$${first.price}` : 'no data yet'}</p> <p>{timing}</p> </div> ); } render(<PriceStream />);
为高频更新跳过 DevTools
使用 WebSocket 或其他实时数据源时,你可能希望跳过将某些高频 action 记录到 DevToolsManager,以免浏览器扩展不堪重负。
import { getDefaultManagers, actionTypes } from '@data-client/react';
import StreamManager from './StreamManager';
import { Ticker } from './Ticker';
export default function getManagers() {
return [
new StreamManager(() => new WebSocket('wss://ws-feed.example.com'), {
ticker: Ticker,
}),
...getDefaultManagers({
devToolsManager: {
// Increase latency buffer for high-frequency updates
latency: 1000,
// Skip WebSocket SET actions to avoid log spam
// (batched writes use the [Ticker] schema)
predicate: (state, action) =>
action.type !== actionTypes.SET ||
(action.schema !== Ticker && action.schema[0] !== Ticker),
},
}),
];
}
加密货币应用
探索 coin-app 示例