所以几天后,我设法找到了两种连接到外部事件总线的方法。我发现我真的不需要外部命令或查询总线,因为它们是通过 API 调用进来的。因此,如果您想使用 NestJS 连接到外部事件总线,这是我发现的两个选项:
- 通过自定义发布者和订阅者
- 通过 NestJS 微服务包
这两种方法的主要区别在于它们连接到外部事件总线的方式以及它们处理传入消息的方式。根据您的需要,一个可能比另一个更适合您,但我选择了第一个选项。
自定义发布者和订阅者
在我的应用程序中,我已经通过使用 NestJS 中的 EventBus 类并为我的事件调用 .publish() 对我的事件总线进行手动发布调用。我创建了一个服务,将本地 NestJS 事件总线与自定义发布者和自定义订阅者一起包裹起来。
eventBusService.ts
export class EventBusService implements IEventBusService {
constructor(
private local: EventBus, // Injected from NestJS CQRS Module
@Inject('eventPublisher') private publisher: IEventPublisher,
@Inject('eventSubscriber') subscriber: IMessageSource) {
subscriber.bridgeEventsTo(this.local.subject$);
}
publish(event: IEvent): void {
this.publisher.publish(event);
};
}
事件服务使用自定义订阅者使用.bridgeEventsTo() 将来自远程事件总线的任何传入事件重定向到本地事件总线。自定义订阅者使用 redis NPM 包的客户端与事件总线通信。
subscriber.ts
export class RedisEventSubscriber implements IMessageSource {
constructor(@Inject('redisClient') private client: RedisClient) {}
bridgeEventsTo<T extends IEvent>(subject: Subject<T>) {
this.client.subscribe('Foo');
this.client.on("message", (channel: string, message: string) => {
const { payload, header } = JSON.parse(message);
const event = Events[header.name];
subject.next(new event(data.event.payload));
});
}
};
此函数还包含将传入的 Redis 事件映射为事件的逻辑。为此,我创建了一个字典,其中包含我在 app.module 中的所有事件,以便我可以查找事件是否知道如何处理传入事件。然后它用一个新事件调用subject.next(),以便将它放在内部的 NestJS 事件总线上。
publisher.ts
为了根据我自己的事件更新其他系统,我创建了一个发布者,将数据发送到 Redis。
export class RedisEventPublisher implements IEventPublisher {
constructor(@Inject('redisClient') private client: RedisClient) {}
publish<T extends IEvent = IEvent>(event: T) {
const name = event.constructor.name;
const request = {
header: {
name
},
payload: {
event
}
}
this.client.publish('Foo', JSON.stringify(request));
}
}
和订阅者一样,这个类使用 NPM 包客户端向 Redis eventBus 发送事件。
微服务设置
某些部分的微服务设置与自定义事件服务方法非常相似。它使用相同的发布者类,但订阅设置的完成方式不同。它使用 NestJS 微服务包设置一个微服务,监听传入的消息,然后调用 eventService 将传入的事件发送到事件总线。
eventService.ts
export class EventBusService implements IEventBusService {
constructor(
private eventBus: EventBus,
@Inject('eventPublisher') private eventPublisher: IEventPublisher,) {
}
public publish<T extends IEvent>(event: T): void {
const data = {
payload: event,
eventName: event.constructor.name
}
this.eventPublisher.publish(data);
};
async handle(string: string) : Promise<void> {
const data = JSON.parse(string);
const event = Events[data.event.eventName];
if (!event) {
console.log(`Could not find corresponding event for
${data.event.eventName}`);
};
await this.eventBus.publish(new event(data.event.payload));
}
}
NestJS 有关于如何设置混合服务的文档,可以在 here 找到。微服务包为您提供了一个 @EventPattern() 装饰器,您可以使用它为传入的事件总线消息创建处理程序,您只需将它们添加到 NestJS 控制器并注入 eventService。
controller.ts
@Controller()
export default class EventController {
constructor(@Inject('eventBusService') private eventBusService:
IEventBusService) {}
@EventPattern(inviteServiceTopic)
handleInviteServiceEvents(data: string) {
this.eventBusService.handle(data)
}
}
因为我不想创建一个混合应用程序来监听传入的事件,所以我决定使用第一个选项。代码很好地组合在一起,而不是使用带有 @EventPattern() 装饰器的随机控制器。
这需要很长时间才能弄清楚,所以我希望它对将来的人有所帮助。 :)