CloudEvents サポート
Spring Integration は CloudEvents 仕様 [GitHub] (英語) をサポートしています。
プロジェクトに次の依存関係を追加します。
ToCloudEventTransformer
ToCloudEventTransformer を使用して、Spring Integration メッセージを CloudEvents 準拠のメッセージに変換します。このトランスフォーマーは CloudEvents 仕様 v1.0 をサポートし、EventFormat または eventFormatContentTypeExpression が指定されている場合は CloudEvents を直列化します。EventFormat または eventFormatContentTypeExpression を指定した場合、トランスフォーマーは EventFormat を使用してペイロードに CloudEvent を生成します。どちらも指定されていない場合、トランスフォーマーはイベントデータメッセージのペイロードをそのまま書き込み、メッセージヘッダーに属性と拡張機能を追加します。トランスフォーマーは式を使用した属性の定義をサポートし、パターンを介してメッセージヘッダー内の拡張機能を識別します。
属性式
SpEL 式を使用して id、source、type、dataSchema、subject の CloudEvents' 属性を設定します。
トランスフォーマーは、CloudEvent インスタンスを作成する時刻に time 属性を設定します。 |
次の表に、デフォルトの式が返す属性名と値を示します。
| 属性名 | デフォルト値 |
|---|---|
| メッセージの ID。 |
| プレフィックス "/spring/" に続いて appName、ピリオド、トランスフォーマーの Bean 名が続きます (例: |
| "spring.message" |
| メッセージの contentType は、デフォルトで |
| 指定されたスキーマへの URI。 |
| イベントプロデューサーのコンテキストにおけるイベントのサブジェクトを識別します。 |
| CloudEvent メッセージが作成された時刻。内部的に現在の時刻に設定されます。この値は変更できませんのでご注意ください。 |
拡張パターン
extensionPatterns コンストラクターパラメーター (文字列の可変長引数) を使用して、ワイルドカード (*) を使用したパターンマッチングを指定します。トランスフォーマーは、任意のパターンに一致するキーを持つメッセージヘッダーを CloudEvent 拡張として含めます。! プレフィックスを使用して、否定によってヘッダーを明示的に除外します。最初に一致したパターン (正負に関わらず) が優先されることに注意してください。
例: パターン "trace*", "span-id", "user-id" を次のように構成します。- trace で始まるヘッダーを含める (例: trace-id、traceparent) - 正確なキー 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 プロセス
変革プロセスを理解する:
CloudEvent Building - CloudEvent 属性を構築します。
Extension Extraction - コンストラクターに渡された extensionPatterns の配列を使用して、CloudEvent 拡張機能を構築します。
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 つの方法のいずれかで設定します。
希望する
EventFormatを設定します。eventFormatContentTypeExpressionには、EventFormatProviderが使用して必要なEventFormatを提供できるコンテンツ型に解決される式を設定します。eventFormatContentTypeExpressionが設定され、コンテンツ型に対応するEventFormatが見つからないためにEventFormatProviderが null を返すと、トランスフォーマーはMessageTransformationExceptionをスローします。eventFormatContentTypeExpressionが解決でき、EventFormatProviderが受け入れるコンテンツ型の例は次のとおりです。application/cloudevents+jsonapplication/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 インスタンスの場合、トランスフォーマーは次のようになります。
CloudEvent データを抽出し、それをメッセージペイロードとして使用します。
CloudEvent 属性(
id、source、type、time、subject、datacontenttype、dataschema)をce-プレフィックスを持つメッセージヘッダーにマッピングします。すべての CloudEvent 拡張機能を、
ce-プレフィックスを持つメッセージヘッダーにマッピングします。ヘッダーキーが
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[] である場合、トランスフォーマーは次のようになります。
content-typeヘッダーを使用して、EventFormatProviderを介して適切なEventFormatを解決します。解決済みのフォーマットを使用して、ペイロードを
CloudEventオブジェクトに逆直列化します。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 属性 | メッセージヘッダーキー | 必須 |
|---|---|---|
|
| はい |
|
| はい |
|
| はい |
|
| いいえ |
|
| いいえ |
|
| いいえ |
|
| いいえ |
extensions |
| いいえ |
出力メッセージの 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 つの方法をサポートしています。
Direct values - 静的値を直接設定する
SpEL 式 - メッセージに基づいて動的な値を取得するには、Spring 式言語を使用します。
関数 - メッセージに基づいて値を評価する 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();
}