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

client から queue に access するには、topic name と connector ID を渡して queue を使用し、QueueManager への reference を作成します。
const queue = squid.queue('topic-name', 'confluent-connector-id');
queue から message を read するには、consume method を使用します。これは、新しい message が queue に post されるたびに update される observable を返します。
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 します。
topicMessagesObs.unsubscribe();
queue に message を追加するには、produce method を使用します。この method は、queue を subscribe している client が consume する message の array を受け取ります。
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 を許可します。
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を参照してください。