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

Confluent Queue

既存の Confluent Cloud Apache Kafka instance を Squid に接続します​

Squid を通じて Confluent Cloud に access するには、まず Squid Console で connector を追加します。

  • Squid Console で Connectors page に移動し、Confluent Kafka queue connector を選択します。

  • 次の detail を入力します。

    • Connector ID - 簡潔で connector の識別に役立つ ID を選択します。
    • Bootstrap servers - comma-separated の server および port list。この value は cluster setting で確認できます。
    • Key - Kafka API Key。cluster をクリックしてから API Keys をクリックし、既存 key を確認するか新しい 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 をクリックします。

Confluent connector

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

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

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

Client code
const topicMessagesObs = queue.consume<MyType>();

topicMessagesObs.subscribe((message: MyType) => {
console.log(message);
});
注記

topic を subscribe すると、受け取る observable は subscription の確立後に produce された message のみを含むよう構成されます。この behavior は、subscription process に server call が含まれ、およそ 100ms の delay が発生するためです。したがって、observable が新しい message の配信を開始するまで短い待機時間が発生します。

memory leak を防ぐため、使用しなくなった observable は complete します。

Client code
topicMessagesObs.unsubscribe();

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

Client code
queue.produce(['hello', 'world']);

Confluent Topic の保護​

client で Confluent 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', 'confluent-connector-id')
allowTopicAccess(): boolean {
return true;
}
}

security function の詳細については、queue security docsを参照してください。