Queue message handler
Queue topic の message を backend service method で処理します。
Queue Message Handler を使用する理由
application が、注文の作成、file の upload、payment の完了などの event を queue に publish し、backend でそれぞれを確実に処理する必要がある場合に使用します。
@onQueueMessage では、method を decorate して deploy するだけです。
- TypeScript
- Python
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
const { orderId, customerEmail, status } = request.message;
// Email the customer an order confirmation using your own email helper.
await this.sendOrderConfirmationEmail(customerEmail, orderId, status);
}
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
event = request['message']
# Email the customer an order confirmation using your own email helper.
await self.send_order_confirmation_email(
event['customerEmail'], event['orderId'], event['status'])
polling も routing logic も不要です。topic に message が到着した瞬間に Squid が handler を呼び出します。
概要
Queue message handler は、指定した topic に message が publish されるたびに自動的に実行される、@onQueueMessage(Python では @on_queue_message)で decorate された backend method です。built-in Squid queue と、Kafka などの external queue integration の両方をサポートします。
Queue message handler を使用する場合
| ユースケース | 推奨 |
|---|---|
| Queue topic に publish された message を処理する | ✅ Queue message handler |
| Database change に反応する | Triggers を使用 |
| Client から function を呼び出す | Executables を使用 |
| Schedule に基づいて code を実行する | Schedulers を使用 |
| External service に HTTP endpoint を公開する | Webhooks を使用 |
仕組み
SquidServiceを拡張する class で、method を@onQueueMessage()で decorate します- topic name と、必要に応じて integration ID を指定します
- Squid が deploy 時に handler を検出して登録します
- topic に message が到着すると、Squid は
QueueMessageRequestobject を指定して method を呼び出します - handler は synchronous にすることも、
Promiseを返すこともできます
クイックスタート
前提条件
- TypeScript
- Python
squid initで初期化された Squid backend project- NPM からインストールされた
@squidcloud/backendpackage
squid initで初期化された Squid backend project- PyPI からインストールされた
squidcloud-backendpackage
ステップ 1: Handler を作成する
SquidService を拡張する service class を作成し、handler method を追加します。
- TypeScript
- Python
import { SquidService, onQueueMessage, QueueMessageRequest } from '@squidcloud/backend';
interface OrderEvent {
orderId: string;
customerEmail: string;
status: 'placed' | 'shipped' | 'delivered';
}
export class OrderService extends SquidService {
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
const { orderId, customerEmail, status } = request.message;
// Email the customer an order confirmation using your own email helper.
await this.sendOrderConfirmationEmail(customerEmail, orderId, status);
}
}
from squidcloud_backend import SquidService, on_queue_message
from squidcloud_backend.types import QueueMessageRequest
from typing import TypedDict
class OrderEvent(TypedDict):
orderId: str
customerEmail: str
status: str # 'placed' | 'shipped' | 'delivered'
class OrderService(SquidService):
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
event: OrderEvent = request['message']
# Email the customer an order confirmation using your own email helper.
await self.send_order_confirmation_email(
event['customerEmail'], event['orderId'], event['status'])
ステップ 2: Service を export する
- TypeScript
- Python
service が service index file から export されていることを確認します。
export * from './example-service';
service が service index file から export されていることを確認します。
from .example_service import OrderService
ステップ 3: Backend を開始または deploy する
ローカル開発では、Squid CLI を使用して backend をローカルで実行します。
squid start
cloud に deploy するには、backend の deployを参照してください。
ステップ 4: 検証する
client または別の service から topic に message を publish し、Squid Console log を確認して handler が呼び出されたことを確認します。
コアコンセプト
Decorator
decorator は 2 つの parameter を受け取ります。
| Parameter | Type | 必須 | 説明 |
|---|---|---|---|
topicName | string | はい | subscribe する queue topic の name |
integrationId | string | いいえ | queue の integration ID。default は built-in Squid queue integration です |
- TypeScript
- Python
// Built-in queue: integrationId defaults to 'built_in_queue'
@onQueueMessage('order-events')
// External integration (e.g. Kafka)
@onQueueMessage('order-events', 'kafka')
# Built-in queue: integration_id defaults to 'built_in_queue'
@on_queue_message('order-events')
# External integration (e.g. Kafka)
@on_queue_message('order-events', 'kafka')
QueueMessageRequest
handler に渡される QueueMessageRequest object には以下が含まれます。
| Property | Type | 説明 |
|---|---|---|
message | T | type 付けされた message payload |
topicName | string | message の送信先 topic の name |
integrationId | string | queue の integration ID |
TypeScript では、generic type parameter T により publisher から handler まで end-to-end の type safety を得られます。Python では QueueMessageRequest は TypedDict のため、dictionary-style lookup で field にアクセスします。
- TypeScript
- Python
interface OrderEvent {
orderId: string;
customerEmail: string;
status: 'placed' | 'shipped' | 'delivered';
}
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
// request.message is typed as OrderEvent
const { orderId, status } = request.message;
}
class OrderEvent(TypedDict):
orderId: str
customerEmail: str
status: str
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
# access fields via dict-style lookups
event: OrderEvent = request['message']
order_id = event['orderId']
topic = request['topicName']
External queue integration の使用
integrationId を省略すると、Squid の built-in queue が使用されます。明示的な integration ID を指定すると、同じ handler pattern を Kafka などの external queue system で使用できます。
- TypeScript
- Python
import { SquidService, onQueueMessage, QueueMessageRequest } from '@squidcloud/backend';
interface OrderEvent {
orderId: string;
customerEmail: string;
status: 'placed' | 'shipped' | 'delivered';
}
export class OrderService extends SquidService {
// Built-in queue
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
console.log('Built-in queue event:', request.message);
}
// External Kafka integration
@onQueueMessage<OrderEvent>('order-events', 'kafka')
async handleKafkaOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
console.log(`Kafka event on topic "${request.topicName}":`, request.message);
}
}
from squidcloud_backend import SquidService, on_queue_message
from squidcloud_backend.types import QueueMessageRequest
class OrderService(SquidService):
# Built-in queue
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
print('Built-in queue event:', request['message'])
# External Kafka integration
@on_queue_message('order-events', 'kafka')
async def handle_kafka_order_event(self, request: QueueMessageRequest) -> None:
print(f"Kafka event on topic '{request['topicName']}':", request['message'])
Error Handling
handler が error を throw すると、その error は Squid Console に log されます。failure を適切に処理するため、logic を try/catch(または try/except)block でラップしてください。
- TypeScript
- Python
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
try {
await this.processOrder(request.message);
} catch (error) {
console.error(`Failed to process order ${request.message.orderId}:`, error);
}
}
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
try:
await self.process_order(request['message'])
except Exception as error:
print(f"Failed to process order {request['message']['orderId']}:", error)
ベストプラクティス
-
Message payload に type を付ける。 TypeScript では
QueueMessageRequest<MyType>generic parameter を、Python ではTypedDictsubclass を使用して、message body に type-safe に access します。 -
Handler 内で error を処理する。 logic を try/catch(または try/except)block でラップし、error を log に記録して、1 つの不正な message が暗黙的に失敗しないようにします。
-
Handler の焦点を絞る。 handler は 1 つのことを行うべきです。message が複数の workflow を trigger する必要がある場合は、すべての logic を handler に記述するのではなく、他の method または service に delegate してください。
-
idempotency を考慮して設計する。 message は、まれに複数回配信されることがあります。duplicate message を処理しても同じ result になるよう handler を設計してください。
コード例
Queue message を external service に転送する
- TypeScript
- Python
import { SquidService, onQueueMessage, QueueMessageRequest } from '@squidcloud/backend';
interface OrderEvent {
orderId: string;
customerEmail: string;
status: 'placed' | 'shipped' | 'delivered';
}
export class OrderService extends SquidService {
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
// Forward each order event to your fulfillment service.
await this.forwardToFulfillment(request.message);
}
}
from squidcloud_backend import SquidService, on_queue_message
from squidcloud_backend.types import QueueMessageRequest
class OrderService(SquidService):
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
# Forward each order event to your fulfillment service.
await self.forward_to_fulfillment(request['message'])
複数の queue integration からの message を処理する
- TypeScript
- Python
import { SquidService, onQueueMessage, QueueMessageRequest } from '@squidcloud/backend';
interface OrderEvent {
orderId: string;
customerEmail: string;
status: 'placed' | 'shipped' | 'delivered';
}
export class OrderService extends SquidService {
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
await this.processOrder(request.message);
}
@onQueueMessage<OrderEvent>('order-events', 'kafka')
async handleKafkaOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
await this.processOrder(request.message);
}
private async processOrder(event: OrderEvent): Promise<void> {
// Email the customer an order confirmation using your own email helper.
await this.sendOrderConfirmationEmail(event.customerEmail, event.orderId, event.status);
}
}
from squidcloud_backend import SquidService, on_queue_message
from squidcloud_backend.types import QueueMessageRequest
from typing import TypedDict
class OrderEvent(TypedDict):
orderId: str
customerEmail: str
status: str
class OrderService(SquidService):
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
await self._process_order(request['message'])
@on_queue_message('order-events', 'kafka')
async def handle_kafka_order_event(self, request: QueueMessageRequest) -> None:
await self._process_order(request['message'])
async def _process_order(self, event: OrderEvent) -> None:
# Email the customer an order confirmation using your own email helper.
await self.send_order_confirmation_email(
event['customerEmail'], event['orderId'], event['status'])
Full-Stack Example
この例では、完全な flow を示します。client が order event を publish し、backend handler がそれを処理します。
Backend: handler を登録する
- TypeScript
- Python
import { SquidService, onQueueMessage, QueueMessageRequest, secureTopic } from '@squidcloud/backend';
interface OrderEvent {
orderId: string;
customerEmail: string;
status: 'placed' | 'shipped' | 'delivered';
}
export class OrderService extends SquidService {
@secureTopic('order-events', 'write')
allowOrderEventPublish(): boolean {
return !!this.getUserAuth();
}
@onQueueMessage<OrderEvent>('order-events')
async handleOrderEvent(request: QueueMessageRequest<OrderEvent>): Promise<void> {
const { orderId, customerEmail, status } = request.message;
console.log(`Received order event: ${orderId} → ${status}`);
// Email the customer an order confirmation using your own email helper.
await this.sendOrderConfirmationEmail(customerEmail, orderId, status);
}
}
from squidcloud_backend import SquidService, on_queue_message, secure_topic
from squidcloud_backend.types import QueueMessageRequest
from typing import TypedDict
class OrderEvent(TypedDict):
orderId: str
customerEmail: str
status: str
class OrderService(SquidService):
@secure_topic('order-events', 'write')
def allow_order_event_publish(self) -> bool:
return self.get_user_auth() is not None
@on_queue_message('order-events')
async def handle_order_event(self, request: QueueMessageRequest) -> None:
event: OrderEvent = request['message']
print(f"Received order event: {event['orderId']} → {event['status']}")
# Email the customer an order confirmation using your own email helper.
await self.send_order_confirmation_email(
event['customerEmail'], event['orderId'], event['status'])
Client: message を publish する
const orderEvent = { orderId: 'order-123', customerEmail: 'jane@example.com', status: 'placed' };
await squid.queue('order-events').produce([orderEvent]);
produce が呼び出されると、Squid は message を topic に配信し、backend handler が自動的に実行されます。
関連項目
- Built-in queue - built-in Squid queue を設定・使用する
- Kafka connector - external Kafka integration に接続する
- Queue security - queue topic への access を保護する
- Triggers - database change に反応する
- Schedulers - schedule に基づいて code を実行する