Kafka Queue
既存の Kafka instance を Squid に接続します
Kafka Connector をセットアップする
built-in queue ではなく独自の Kafka connector を使用する場合は、まず connector を Squid application に追加します。
-
Squid Console で Connectors page に移動し、Kafka connector を選択します。
-
次の detail を指定します。
- Connector ID - 簡潔で connector の識別に役立つ ID を選択します。
- Bootstrap servers - server と port の comma-separated list。
- Key - Kafka API Key。
- Secret - Kafka API secret。Squid Secrets に安全に保存されます。
- Avro schema registry - Avro message を consume または produce する場合、Avro schema registry の URL を指定します。この value は cluster setting で確認できます。
-
Test connection をクリックして connector information を検証します。検証後、Add connector をクリックします。

client から queue に access するには、topic name と connector ID を渡して queue を使用し、QueueManager への reference を作成します。
const queue = squid.queue('topic-name', 'kafka-connector-id');
queue から message を read するには、consume method を使用します。これは、新しい message が queue に post されるたびに update される observable を返します。message は string type です。
const topicMessagesObs = queue.consume();
topicMessagesObs.subscribe((message: string) => {
console.log(message);
});
topic を subscribe すると、受け取る observable は subscription 確立後に produce された message のみを含むように構成されます。この動作は、subscription process に server call が含まれ、およそ 100ms の delay が発生するためです。したがって、observable が新しい message の配信を開始するまでに短い待機時間が発生します。
memory leak を防ぐため、使用しなくなった observable を complete します。
topicMessagesObs.unsubscribe();
queue に message を追加するには、produce method を使用します。この method は string message の array を受け取り、queue を subscribe している client により consume されます。
queue.produce(['hello', 'world']);
Apache Kafka Topic の保護
client で CApache Kafka topic に access するには、security function が必要です。
Squid queue topic を保護するには、topic name と action type を渡して @secureTopic decorator を使用します。次の code は、'topic-name' topic の queue に read と write access を許可します。
import { SquidService, secureTopic } from '@squidcloud/backend';
export class ExampleService extends SquidService {
@secureTopic('topic-name', 'all', 'kafka-connector-id')
allowTopicAccess(): boolean {
return true;
}
}
security function の詳細については、queue security docsを参照してください。