跳到主要内容
备注

这篇翻译之后英文原文有更新,内容可能已过时。

Manager 与 Middleware

Reactive Data Client 采用 flux store 模式,其特点是 store 的单向数据流易于理解和调试。状态更新由 reducer 函数执行。

Manager 的 flux 流程Manager 的 flux 流程

在 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);
}
}
index.tsx
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 />);
结果
Store▶

为高频更新跳过 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),
},
}),
];
}

加密货币应用​

更多演示