Vello: クラスパスにない外部イベント処理

Vello: クラスパスにない外部イベント処理

問題の定義

Lakeyプロジェクトでは、社内メッセージングフレームワークであるPrologueのVelloモジュールを使用して外部イベントを保存し、LLMコンテキストで利用できるよう、外部イベントを管理しています。

Prologueがデフォルトで提供するイベントコンバーターである'TypeAwareEventConverter'は、メッセージを受信するとヘッダーの'payloadClass'値を読み取り、クラス名が稼働中のサービスのクラスパスに存在するかを確認した後、デシリアライズを実行します。

この方式は、同一サービス内、または外部サービスのEventモジュールに依存した状態でイベントを送受信する場合には問題ありません。すでにEventクラスがサービスのクラスパスに存在するためです。しかし、未定義のまま発行されたイベントの場合は事情が異なります。未定義イベントのクラスはLakeyのクラスパスに存在しないため例外が発生し、最終的にそのイベントは正常に処理されないまま失われます。

外部サービスのイベントスキーマが変更されるたびにLakey内部のロジックも変更されるため、サービス間に強い結合が生じます。外部サービスのEventモジュールに依存せずにイベントを保存する方法を検討し、この問題を解決するために'LakeyEventConverter'を直接実装しました。

既存のEventConverterの動作方式

1. メッセージペイロードを指定されたエンコーディングで文字列に変換します。

2. ヘッダーの'payloadClass'値が空の場合は、'Map'型としてデシリアライズして返します。

3. 'payloadClass'値が存在する場合は、'Class.forName()'でクラスをロードした後、その型としてデシリアライズします。

4. クラスのロードに失敗すると、最終的に'EventConversionFailed'例外が発生します。

解決方法:LakeyEventConverter

4番目の問題を解決するため、'TypeAwareEventConverter'を継承したうえで'convert()'メソッドをオーバーライドしました。型変換に失敗したイベントを破棄せず、解釈可能な原始的な形式で保持する方法です。

public class LakeyEventConverter extends TypeAwareEventConverter {

    private final String encoding;

    public LakeyEventConverter(
            @Value("${vizend.prologue.vello.encoding:UTF-8}") String encoding) {
        super(encoding);
        this.encoding = encoding;
    }

    @Override
    public Object convert(VelloMessage message) {
        try {
            Object object = super.convert(message);
            if (object instanceof LakeyOutgoingDomainEvent out) {
                return out.toIncomingEvent();
            }
            return object;
        } catch (EventConvertException e) {
            String payloadJson = new String(message.getPayload(), Charset.forName(encoding));
            Map<String, Object> payload = JsonUtil.fromJson(payloadJson, Map.class);

            String fqcn = message.getPayloadClass();
            String simpleClassName = fqcn.substring(fqcn.lastIndexOf('.') + 1);

            return UnreferencedEvent.builder()
                    .eventType(simpleClassName)
                    .headers(message.getHeaders())
                    .payload(payload)
                    .messageId(message.getId())
                    .subject(message.getSubject())
                    .receivedAt(LocalDateTime.now())
                    .build();
        }
    }
}

'LakeyEventConverter'の処理フローは、2つのケースに分かれます。

1. 通常どおり変換に成功した場合

親コンバーター('TypeAwareEventConverter')の変換が成功し、Lakey独自のイベントまたは依存しているイベントクラスである場合は、変換されたオブジェクトをそのまま返します。この場合、別途後処理を行う必要はありません。

2. 外部イベントのように型変換に失敗した場合

親クラスが'EventConvertException'をスローした場合は、これを捕捉して別途処理します。この場合、ペイロードを'Map<String, Object>'としてデシリアライズし、ヘッダーに含まれるFQCNから単純なクラス名だけを抽出して'eventType'として使用します。その後、これを基に'UnreferencedEvent'オブジェクトを生成して返します。

このとき、単にペイロードだけを残すのではなく、以下の情報も併せて保持します。

* メッセージID、受信トピック(subject)、メッセージヘッダー、受信時刻、元のペイロード

UnreferencedEvent

UnreferencedEventは、参照するクラスが存在しないイベントを格納するオブジェクトです。名前のとおり、Lakeyは型を把握できないものの、収集する価値のあるイベントを表します。

@Getter
@Builder
public class UnreferencedEvent {
    private String eventType;
    private Map<String, String> headers;
    private Map<String, Object> payload;
    private LocalDateTime receivedAt;
    private String messageId;
    private String subject;
}

ここで重要な値はペイロードだけではありません。messageIdは重複処理のために必要で、subjectはイベントが到着したトピックを確認するために使用されます。headersとreceivedAtは、後からイベントの出所と受信時点を追跡するのに役立ちます。

UnreferencedEventからRawEventまで

LakeyEventConverterがUnreferencedEventを返すと、Prologueの@EventHandlerが付与されたUnreferencedEventHandlerがこのイベントを受け取ります。ハンドラーは特別なビジネス上の判断を行わず、RawEventFlow.register()に渡します。

@Component
@RequiredArgsConstructor
public class UnreferencedEventHandler {

    private final RawEventFlow rawEventFlow;

    @EventHandler
    public void handleUnreferencedEvent(UnreferencedEvent unreferencedEvent) {
        if (unreferencedEvent == null) {
            log.warn("UnreferencedEvent is null");
            return;
        }
        rawEventFlow.register(unreferencedEvent);
    }
}

RawEventFlow.register()では、3つの項目を確認します。

public void register(UnreferencedEvent unreferencedEvent) {
    if (!rawEventAction.isTargetEvent(unreferencedEvent)) {
        return;
    }

    RawEventCdo rawEventCdo = rawEventAction.convertToRawEventCdo(unreferencedEvent);

    if (rawEventLogic.existsByEventId(rawEventCdo.getEventId())) {
        return;
    }

    String rawEventId = rawEventLogic.registerRawEvent(rawEventCdo);
    eventProxy.produceEvent(RawEventProjectionRequestedEvent.newInstance(rawEventId));
}

1. 対象トピックかどうかを確認

外部から受信したすべてのイベントを保存するわけではありません。LakeyConfigPropertiesに設定されたtargetTopicsに含まれるトピックだけを処理します。収集対象ではないイベントを早い段階で除外することで、不要な保存や後続処理を削減します。

2. 重複イベントかどうかを確認

messageIdをeventIdとして使用し、すでに保存されたイベントかどうかを確認します。同じイベントが再度届いた場合でも、1回だけ保存されるようにします。

3. RawEventとして保存

対象トピックであり、重複でもない場合は、UnreferencedEventをRawEventCdoに変換して保存します。保存後にはRawEventProjectionRequestedEventを発行します。実際のデータ前処理は、このイベントを処理するハンドラーの段階で行われます。

まとめ

この構造を作った理由は明確です。イベントの型を事前に把握できなくても、受信したイベント自体を失わないためです。

従来の'TypeAwareEventConverter'は、クラスパスに存在しない型に遭遇すると例外をスローして処理を中断していました。一方、'LakeyEventConverter'は、例外処理されたイベントを'UnreferencedEvent'というCommon Entityに格納して保持します。その結果、Lakeyは外部サービスのイベントクラスに直接依存しなくても、イベントを収集できるようになりました。

イベント処理においてサービス間の結合度を下げる必要がある場合に、参考にしていただければ幸いです。

Ted

Site footer