CloudEvents サポート

Spring Integration は CloudEvents 仕様 [GitHub] (英語) をサポートしています。

プロジェクトに次の依存関係を追加します。

<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-cloudevents</artifactId>
    <version>7.1.1</version>
</dependency>
implementation "org.springframework.integration:spring-integration-cloudevents:7.1.1"

ToCloudEventTransformer

ToCloudEventTransformer を使用して、Spring Integration メッセージを CloudEvents 準拠のメッセージに変換します。このトランスフォーマーは CloudEvents 仕様 v1.0 をサポートし、EventFormat または eventFormatContentTypeExpression が指定されている場合は CloudEvents を直列化します。EventFormat または eventFormatContentTypeExpression を指定した場合、トランスフォーマーは EventFormat を使用してペイロードに CloudEvent を生成します。どちらも指定されていない場合、トランスフォーマーはイベントデータメッセージのペイロードをそのまま書き込み、メッセージヘッダーに属性と拡張機能を追加します。トランスフォーマーは式を使用した属性の定義をサポートし、パターンを介してメッセージヘッダー内の拡張機能を識別します。

属性式

SpEL 式を使用して idsourcetypedataSchemasubject の CloudEvents' 属性を設定します。

トランスフォーマーは、CloudEvent インスタンスを作成する時刻に time 属性を設定します。

次の表に、デフォルトの式が返す属性名と値を示します。

属性名 デフォルト値

id

メッセージの ID。

source

プレフィックス "/spring/" に続いて appName、ピリオド、トランスフォーマーの Bean 名が続きます (例: /spring/myapp.toCloudEventTransformerBean)。

type

"spring.message"

dataContentType

メッセージの contentType は、デフォルトで application/octet-stream になります。その他の例としては、application/jsonapplication/x-avroapplication/xml などがあります。

dataSchema

指定されたスキーマへの URI。dataSchema のデフォルトは null です。

subject

イベントプロデューサーのコンテキストにおけるイベントのサブジェクトを識別します。subject のデフォルトは null です。

time

CloudEvent メッセージが作成された時刻。内部的に現在の時刻に設定されます。この値は変更できませんのでご注意ください。

拡張パターン

extensionPatterns コンストラクターパラメーター (文字列の可変長引数) を使用して、ワイルドカード (*) を使用したパターンマッチングを指定します。トランスフォーマーは、任意のパターンに一致するキーを持つメッセージヘッダーを CloudEvent 拡張として含めます。! プレフィックスを使用して、否定によってヘッダーを明示的に除外します。最初に一致したパターン (正負に関わらず) が優先されることに注意してください。

例: パターン "trace*", "span-id", "user-id" を次のように構成します。- trace で始まるヘッダーを含める (例: trace-idtraceparent) - 正確なキー span-id と user-id を持つヘッダーを含める - 一致するすべてのヘッダーを CloudEvent の拡張として追加する

特定のヘッダーを除外するには、否定パターンを使用します。"custom-*", "!custom-internal" は、custom-internal を除く custom- で始まるすべてのヘッダーを含みます。

DSL を使用した設定

CloudEvents ファクトリを使用して、Java DSL を使用するフローに ToCloudEventTransformer を追加します。

@Bean
public ToCloudEventTransformer cloudEventTransformer() {
    return new ToCloudEventTransformer("trace*", "correlation-id");
}

@Bean
public IntegrationFlow cloudEventTransformFlow(ToCloudEventTransformer toCloudEventTransformer) {
    return IntegrationFlows
        .from("inputChannel")
        .transform(CloudEvents.toCloudEventTransformer().get())
        .channel("outputChannel")
        .get();
}

CloudEventTransformer プロセス

変革プロセスを理解する:

  1. CloudEvent Building - CloudEvent 属性を構築します。

  2. Extension Extraction - コンストラクターに渡された extensionPatterns の配列を使用して、CloudEvent 拡張機能を構築します。

  3. Format Conversion - 指定された EventFormat を適用するか、設定されていない場合はバイナリフォーマットモードで変換を処理します。

基本的な変換は、次のようなパターンを持つ場合があります。

// Input message with headers
Message<byte[]> inputMessage = MessageBuilder
    .withPayload("Hello CloudEvents".getBytes(StandardCharsets.UTF_8))
    .withHeader(MessageHeaders.CONTENT_TYPE, "text/plain")
    .build();

ToCloudEventTransformer transformer = new ToCloudEventTransformer();

// Transform to CloudEvent
Object cloudEventMessage = transformer.transform(inputMessage);

EventFormats

ToCloudEventTransformer は、EventFormat が利用可能な場合はフォーマットを使用して CloudEvent をメッセージのペイロードに直列化し、それ以外の場合はバイナリフォーマットモードを使用します。EventFormat は、次の 2 つの方法のいずれかで設定します。

  1. 希望する EventFormat を設定します。

  2. eventFormatContentTypeExpression には、EventFormatProvider が使用して必要な EventFormat を提供できるコンテンツ型に解決される式を設定します。eventFormatContentTypeExpression が設定され、コンテンツ型に対応する EventFormat が見つからないために EventFormatProvider が null を返すと、トランスフォーマーは MessageTransformationException をスローします。eventFormatContentTypeExpression が解決でき、EventFormatProvider が受け入れるコンテンツ型の例は次のとおりです。

    • application/cloudevents+json

    • application/cloudevents+xml

EventFormat と eventFormatContentTypeExpression が設定されていない場合、トランスフォーマーはクラウドイベントプレフィックス(デフォルトは ce-)を使用してメッセージヘッダーにクラウドイベント属性と拡張機能を追加し、ペイロードは変更しません(バイナリフォーマットモード)。

特定の EventFormat を使用するには、関連する依存関係を追加します。例: XML EventFormat を追加するには、次の依存関係 io.cloudevents:cloudevents-xml を追加します。使用可能なイベント形式については、CloudEvents Java リファレンスドキュメント (英語) を参照してください。

CloudEvents に変換するメッセージのペイロードが byte[] 型であることを確認してください。ペイロードがバイト配列でない場合、トランスフォーマーは IllegalArgumentException 例外をスローします。

FromCloudEventTransformer

FromCloudEventTransformer を使用して、CloudEvents メッセージを Spring Integration メッセージに変換します。このトランスフォーマーは、CloudEvents 仕様 v1.0 をサポートし、CloudEvent オブジェクトまたは直列化された CloudEvent バイト配列の 2 種類のペイロード型から CloudEvents を処理します。

トランスフォーマーは、メッセージペイロードから CloudEvent データを抽出し、CloudEvent 属性と CloudEvent 拡張機能を ce- プレフィックス付きのメッセージヘッダーにマッピングします。

サポートされているペイロードの種類

トランスフォーマーは、以下のペイロード型のメッセージを受け入れます。

CloudEvent オブジェクト型

メッセージペイロードが CloudEvent インスタンスの場合、トランスフォーマーは次のようになります。

  1. CloudEvent データを抽出し、それをメッセージペイロードとして使用します。

  2. CloudEvent 属性(idsourcetypetimesubjectdatacontenttypedataschema)を ce- プレフィックスを持つメッセージヘッダーにマッピングします。

  3. すべての CloudEvent 拡張機能を、ce- プレフィックスを持つメッセージヘッダーにマッピングします。

  4. ヘッダーキーが CloudEvent 属性または拡張機能と一致する場合を除き、すべての元のメッセージヘッダーが保持されます。一致する場合は、元の値が上書きされます。

例:

String orderJson = ...

CloudEvent cloudEvent = CloudEventBuilder.v1()
    .withId("event-123")
    .withSource(URI.create("/myapp/orders"))
    .withType("order.created")
    .withData("application/json", orderJson.getBytes())
    .withExtension("traceid", "trace-abc")
    .build();

Message<CloudEvent> inputMessage = MessageBuilder
    .withPayload(cloudEvent)
    .build();

FromCloudEventTransformer transformer = new FromCloudEventTransformer();
Message<?> outputMessage = transformer.transform(inputMessage);

上記の例の outputMessage は、以下のような出力を生成します。

GenericMessage [
    payload = byte[13],
    headers = {
        ce-source          = /myapp/orders,
        ce-datacontenttype = application/json,
        ce-type            = order.created,
        ce-id              = event-123,
        ce-traceid         = trace-abc,
        id                 = 2df76f27-d139-424c-19b6-80b64e4a33b0,
        contentType        = application/json,
        timestamp          = 1770667476433
    }
]

シリアル番号付き CloudEvent 型

メッセージペイロードが直列化された CloudEvent を含む byte[] である場合、トランスフォーマーは次のようになります。

  1. content-type ヘッダーを使用して、EventFormatProvider を介して適切な EventFormat を解決します。

  2. 解決済みのフォーマットを使用して、ペイロードを CloudEvent オブジェクトに逆直列化します。

  3. CloudEvent オブジェクト型の項に記載されている手順と同じ手順に従います。

サポートされているコンテンツ型に関する情報は、EventFormats のセクションで説明されています。

FromCloudEventTransformer を使用すると、EventFormatProvider が contentType ヘッダーに対応する EventFormat を見つけられなかった場合、またはメッセージに contentType ヘッダーが含まれていない場合に使用される EventFormat を設定できます。設定されていない場合、EventFormatProvider が EventFormat を見つけられなかった場合は、MessageTransformationException がスローされます。

例:

byte[] serializedCloudEvent = """
				{
					"specversion": "1.0",
					"id": "316b0cf3-0c4d-5858-6bd2-863a2042f442",
					"source": "/spring/testapp.jsonTransformerWithExtensions",
					"type": "spring.message",
					"subject": "test.subject",
					"datacontenttype": "text/plain",
					"time": "2026-01-30T08:53:06.099486-05:00",
					"traceid": "trace-123",
					"data": "Hello, World!"
				}
				""";

Message<byte[]> inputMessage = MessageBuilder
    .withPayload(serializedCloudEvent)
    .setHeader(MessageHeaders.CONTENT_TYPE, "application/cloudevents+json")
    .build();

FromCloudEventTransformer transformer = new FromCloudEventTransformer();
Message<?> outputMessage = transformer.transform(inputMessage);

上記の例の outputMessage は、以下のような出力を生成します。

GenericMessage [
    payload = byte[13],
    headers = {
        ce-source           = /spring/testapp.jsonTransformerWithExtensions,
        ce-datacontenttype  = text/plain,
        ce-subject          = test.subject,
        ce-type             = spring.message,
        ce-id               = 316b0cf3-0c4d-5858-6bd2-863a2042f442,
        ce-traceid          = trace-123,
        ce-time             = 2026-01-30T08:53:06.099486-05:00,
        id                  = 463c0878-a9cb-7269-a503-b4224088cd42,
        contentType         = text/plain,
        timestamp           = 1770392214225
    }
]

CloudEvent 属性マッピング

トランスフォーマーは、以下の CloudEventHeaders 定数を使用して、CloudEvent 属性をメッセージヘッダーにマッピングします。

CloudEvent 属性 メッセージヘッダーキー 必須

id

ce-id

はい

source

ce-source

はい

type

ce-type

はい

time

ce-time

いいえ

subject

ce-subject

いいえ

datacontenttype

ce-datacontenttype

いいえ

dataschema

ce-dataschema

いいえ

extensions

ce-{extensionName}

いいえ

出力メッセージの contentType ヘッダーは、常に CloudEvent の datacontenttype 値に設定されます。

DSL を使用した設定

CloudEvents ファクトリを使用して、Java DSL を使用するフローに FromCloudEventTransformer を追加します。

@Bean
public FromCloudEventTransformer fromCloudEventTransformer() {
    return new FromCloudEventTransformer();
}

@Bean
public IntegrationFlow fromCloudEventFlow(FromCloudEventTransformer fromCloudEventTransformer) {
    return IntegrationFlows
        .from("cloudEventInputChannel")
        .transform(CloudEvents.fromCloudEventTransformer())
        .channel("messageOutputChannel")
        .get();
}

CloudEvent ヘッダービルダー

CloudEventHeadersBuilder は、Java DSL を使用する際に CloudEvent 属性を追加するために使用できます。ビルダーは、CloudEvents.headers() ファクトリメソッドを介してアクセスできます。

ビルダーは、ヘッダー値を設定するための 3 つの方法をサポートしています。

  1. Direct values - 静的値を直接設定する

  2. SpEL 式 - メッセージに基づいて動的な値を取得するには、Spring 式言語を使用します。

  3. 関数 - メッセージに基づいて値を評価する Java 関数を提供する

基本的な使い方

@Bean
public IntegrationFlow cloudEventFlow() {
    return IntegrationFlow
        .from("inputChannel")
        .enrichHeaders(CloudEvents.headers()
            .idExpression("headers.orderId")
            .sourceFunction(msg -> URI.create("https://example.com/" + msg.getHeaders().get("action")))
            .type("order.created")
            .subject("order-processing"))
        .transform(CloudEvents.toCloudEventTransformer())
        .channel("outputChannel")
        .get();
}