All Products
Search
Document Center

ApsaraMQ for Kafka:Buat konektor Elasticsearch sink

Last Updated:Jun 21, 2026

Konektor ini mengekspor data dari topik sumber di instans Message Queue for Apache Kafka ke Alibaba Cloud Elasticsearch.

Prasyarat

Untuk informasi selengkapnya, lihat Prasyarat.

Langkah 1: Buat sumber daya layanan target

Langkah 2: Buat konektor Elasticsearch sink

  1. Masuk ke ApsaraMQ for Kafka console. Pada halaman Overview, pilih wilayah di bagian Resource Distribution.

  2. Di panel navigasi sebelah kiri, pilih Connector Ecosystem Integration > Tasks.

  3. Pada halaman Tasks, klik Create Task.

  4. Pada halaman Create Task, konfigurasikan parameter Task Name dan Description. Lalu, ikuti petunjuk di layar untuk mengonfigurasi parameter lainnya.

    • Konfigurasi tugas

      1. Pada wizard konfigurasi Source, atur Data Provider ke ApsaraMQ for Kafka, konfigurasikan parameter berikut, lalu klik Next.

        Parameter

        Description

        Example

        Region

        Wilayah tempat instans ApsaraMQ for Kafka berada.

        China (Hangzhou)

        ApsaraMQ for Kafka Instance

        ID instans ApsaraMQ for Kafka tempat data yang ingin Anda rutekan diproduksi.

        alikafka_post-cn-9hdsbdhd****

        Topic

        Topik pada instans ApsaraMQ for Kafka tempat data yang ingin Anda rutekan diproduksi.

        guide-sink-topic

        Group ID

        ID grup pada instans ApsaraMQ for Kafka tempat data yang ingin Anda rutekan diproduksi.

        • Quickly Create: Sistem secara otomatis membuat grup dengan ID dalam format GID_EVENTBRIDGE_xxx.

        • Use Existing Group: Pilih ID grup yang sudah ada dan tidak sedang digunakan. Jika Anda memilih grup yang sedang digunakan, publikasi dan langganan pesan yang ada akan terpengaruh.

        Use Existing Group

        Consumer Offset

        • Latest Offset: Pesan dikonsumsi dari offset terbaru.

        • Earliest Offset: Pesan dikonsumsi dari offset paling awal.

        Latest Offset

        Network Configuration

        Jika diperlukan transmisi data lintas batas negara, pilih Self-managed Internet. Dalam kasus lain, pilih Basic Network.

        Basic Network

        Data Format

        Fitur format data digunakan untuk mengencode data biner yang dikirimkan dari sumber ke format data tertentu. Beberapa format data didukung. Jika Anda tidak memiliki persyaratan khusus terkait encoding, tentukan Json sebagai nilainya.

        • Json: Mengencode data biner ke format JSON berdasarkan UTF-8 dan menempatkannya di muatan (payload).

        • Text: Format default. Mengencode data biner ke string berdasarkan UTF-8 dan menempatkannya di muatan (payload).

        • Binary: Data biner diencode ke string berdasarkan encoding Base64 lalu dimasukkan ke muatan (payload).

        Text

        Messages

        Parameter Advanced Configuration. Jumlah maksimum pesan yang dikirim dalam satu batch. Permintaan hanya dikirim ketika jumlah pesan yang terakumulasi mencapai nilai ini. Nilainya harus berupa bilangan bulat antara 1 hingga 10.000.

        2000

        Interval (Unit: Seconds)

        Parameter Advanced Configuration. Interval pemanggilan fungsi. Sistem mengagregasi pesan dan mengirimkannya ke Function Compute pada interval ini. Nilainya harus berupa bilangan bulat antara 0 hingga 15. Satuannya adalah detik. Nilai 0 berarti pesan dikirimkan segera.

        3

      2. Pada langkah Filtering, definisikan pola event untuk menyaring permintaan. Untuk informasi selengkapnya, lihat Event patterns.

      3. Pada langkah Transformation, konfigurasikan pembersihan data untuk pemrosesan kompleks, seperti pemisahan, pemetaan, pengayaan, dan perutean dinamis. Untuk informasi selengkapnya, lihat Gunakan Function Compute untuk membersihkan data pesan.

      4. Pada langkah Sink, pilih Alibaba Cloud Elasticsearch acs.elasticSearch untuk Service Type dan konfigurasikan parameter berikut.

        Parameter

        Description

        Example

        Elasticsearch Cluster

        Instans Elasticsearch yang telah Anda buat.

        es-cn-pe336j0gj001e****

        Cluster Logon Name

        Nama logon untuk instans, yaitu elastic secara default.

        elastic

        Instance logon password

        Password yang Anda konfigurasi saat membuat instans.

        ******

        Index Name

        Nama indeks yang telah Anda buat. Untuk informasi lebih lanjut tentang cara membuat indeks, lihat Memulai. Nama indeks dapat berupa konstanta string atau variabel JSONPath, seperti product_info atau $.data.key.

        product_info

        Document Type

        Jenis dokumen data. Ini dapat berupa konstanta string atau variabel yang diekstraksi menggunakan ekspresi JSONPath.

        Contoh: _doc atau $.data.key.

        Catatan

        Parameter ini hanya dapat dikonfigurasi untuk versi instans Elasticsearch sebelum 7.0. Nilai default-nya adalah konstanta _doc.

        _doc

        Document

        Pilih apakah akan mengirimkan event lengkap atau sebagian ke Elasticsearch. Untuk event sebagian, Anda harus mengonfigurasi aturan ekstraksi JSONPath.

        Complete Event

        Network Configuration

        • VPC: Kirim pesan Kafka ke Elasticsearch melalui VPC.

        • Public Network: Kirim pesan Kafka ke Elasticsearch melalui Internet.

        Public Network

        VPC

        VPC yang berisi instans Elasticsearch. Parameter ini hanya diperlukan jika Network Configuration diatur ke VPC.

        vpc-bp17fapfdj0dwzjkd****

        vSwitch

        vSwitch tempat instans Elasticsearch berada. Parameter ini hanya diperlukan jika Network Configuration diatur ke VPC.

        vsw-bp1gbjhj53hdjdkg****

        Security Group

        Grup keamanan. Parameter ini hanya diperlukan jika Network Configuration diatur ke VPC.

        test_group

    • Properti tugas

      Konfigurasikan kebijakan pengulangan (retry policy) untuk pengiriman event yang gagal dan metode penanganan error. Untuk informasi selengkapnya, lihat Retries and dead-letter queues.

  5. Setelah selesai mengonfigurasi, klik Save. Di halaman Tasks, temukan tugas konektor Elasticsearch sink yang telah Anda buat. Kolom Status menampilkan Starting. Saat status berubah menjadi Running, konektor telah dibuat dan siap digunakan.

Langkah 3: Uji konektor Elasticsearch sink

  1. Di halaman Tasks, klik topik sumber di kolom Event Source pada tugas konektor Elasticsearch sink.

  2. Di halaman detail topik, klik Send Test Message.
  3. Di panel Start to Send and Consume Message, konfigurasikan pesan sebagai berikut, lalu klik OK.

    Atur metode pengiriman ke Console. Atur Message Key ke es-sink-k1 dan Message Content ke {"esk1":1,"esk2":"v2"}. Untuk Send to Specified Partition, pilih No.

  4. Masuk ke Konsol Elasticsearch dan akses instans melalui Kibana. Untuk informasi selengkapnya, lihat Memulai.

  5. Di konsol Kibana, jalankan perintah berikut untuk melihat hasil penyisipan data.

    GET /your-index-name/_search

    Hasil penyisipan data: Kueri mengembalikan status 200 OK dan mencakup satu dokumen, di mana _index adalah product_info dan _id adalah 1717558528. Bidang _source berisi bidang topic, partition, offset, timestamp, headers, key, dan value. key adalah es-sink-k1 dan value adalah {"esk1": 1, "esk2": "v2"}. Hal ini mengonfirmasi bahwa data berhasil ditulis ke Elasticsearch.