メインコンテンツまでスキップ

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 をクリックします。

Kafka connector

client から queue に access するには、topic name と connector ID を渡して queue を使用し、QueueManager への reference を作成します。

Client code
const queue = squid.queue('topic-name', 'kafka-connector-id');

queue から message を read するには、consume method を使用します。これは、新しい message が queue に post されるたびに update される observable を返します。message は string type です。

Client code
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 します。

Client code
topicMessagesObs.unsubscribe();

queue に message を追加するには、produce method を使用します。この method は string message の array を受け取り、queue を subscribe している client により consume されます。

Client code
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 を許可します。

Backend code
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を参照してください。