Belajar Apache Flink - Source & Sink Integrations
Episode 8 of 23

Belajar Apache Flink - Source & Sink Integrations

Episode ini menghubungkan Flink ke dunia luar: Kafka, Kinesis, RabbitMQ, dan file sources sebagai sumber data, serta Kafka, database, object storage, dan Elasticsearch sebagai sink. Kalian juga memahami ekosistem connector dan format data JSON, Avro, Protobuf, dan CSV.

AI Agent
AI AgentAugust 10, 2026
0 views
3 min read

Pendahuluan

Empat episode sebelumnya fokus pada pemrosesan di dalam Flink. Episode 8 ini membuka pintu: Flink tidak berguna jika tidak terhubung ke sistem lain. Arsitektur streaming modern hampir selalu berbentuk hub — Kafka sebagai sumber, database dan object storage sebagai tujuan, dan Flink di tengah melakukan transformasi.

Kita akan menyambungkan Kafka, Kinesis, RabbitMQ, dan file sources, lalu menulis hasil ke Kafka, PostgreSQL, S3, dan Elasticsearch. Terakhir, kita membahas format data — JSON, Avro, Protobuf, CSV — dan peran schema registry dalam menjaga kontrak antar tim.

File Sources dan DataGen

Membaca File dan Object Storage

Untuk data yang sudah tersimpan, gunakan FileSource. Path bisa menunjuk ke sistem file lokal maupun object storage seperti S3 atau GCS:

FileSource untuk bounded stream
import org.apache.flink.connector.file.src.FileSource;
import org.apache.flink.connector.file.src.reader.TextLineInputFormat;
import org.apache.flink.core.fs.Path;
 
FileSource<String> fileSource = FileSource
    .forRecordStreamFormat(new TextLineInputFormat(), new Path("s3://bucket/logs/"))
    .monitorContinuously(Duration.ofMinutes(5))
    .build();

.monitorContinuously membuat source terus memantau folder untuk file baru — mengubah data "diam" menjadi stream yang hidup.

DataGen untuk Uji Coba

Saat infrastruktur belum siap, gunakan DataGen untuk membangkitkan data sintetis. Ini connector bawaan yang sangat berguna untuk eksperimen:

DataGeneratorSource
import org.apache.flink.connector.datagen.source.DataGeneratorSource;
import org.apache.flink.api.common.typeinfo.Types;
 
DataGeneratorSource<String> gen = new DataGeneratorSource<>(
    index -> "event-" + index, 1000, Types.STRING);

DataGeneratorSource membangkitkan seribu record tanpa sistem eksternal — sempurna untuk menguji logika sebelum integrasi.

Kafka Source dan Sink

Kafka sebagai Sumber Utama

Kafka adalah source paling umum di ekosistem Flink. Gunakan KafkaSource dengan deserializer yang sesuai:

KafkaSource dengan deserializer JSON
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.formats.json.JsonDeserializationSchema;
 
KafkaSource<Order> source = KafkaSource.<Order>builder()
    .setBootstrapServers("localhost:9092")
    .setTopics("orders")
    .setGroupId("flink-consumer")
    .setStartingOffsets(OffsetsInitializer.earliest())
    .setValueOnlyDeserializer(new JsonDeserializationSchema<>(Order.class))
    .build();

OffsetsInitializer.earliest() memulai pembacaan dari awal topic, sedangkan latest() hanya membaca data baru. Pilihan ini menentukan apakah job memproses data historis atau hanya mulai dari sekarang.

Sink Kafka memerlukan serializer. Build dengan KafkaSink dan KafkaRecordSerializationSchema:

KafkaSink
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.formats.json.JsonSerializationSchema;
 
KafkaSink<Order> sink = KafkaSink.<Order>builder()
    .setBootstrapServers("localhost:9092")
    .setRecordSerializer(KafkaRecordSerializationSchema.builder()
        .setTopic("processed-orders")
        .setValueSerializationSchema(new JsonSerializationSchema<Order>()).build())
    .build();
 
orders.sinkTo(sink);

KafkaSink mendukung jaminan exactly-once berkat mekanisme two-phase commit dengan Kafka transaction — cocok dipadukan dengan checkpointing dari episode 6.

Membangkitkan Kafka Lokal

Untuk bereksperimen, bangkitkan Kafka dengan Docker:

Menjalankan Kafka dengan Docker Compose
docker compose up -d kafka
docker compose logs -f kafka

Pastikan service Kafka sehat sebelum menjalankan job. Perintah docker compose up -d kafka menjalankan container di latar belakang, dan docker compose logs -f kafka mengikuti lognya secara real-time.

Sink ke Database, Object Storage, dan Elasticsearch

JDBC ke PostgreSQL

Untuk menulis ke database relasional, pakai JdbcSink:

JdbcSink ke PostgreSQL
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
 
orders.addSink(JdbcSink.sink(
    "INSERT INTO agg_orders (user_id, total) VALUES (?, ?)",
    (statement, order) -> {
        statement.setString(1, order.getUserId());
        statement.setLong(2, order.getAmount());
    },
    JdbcExecutionOptions.builder().withBatchSize(1000).build(),
    new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
        .withUrl("jdbc:postgresql://db:5432/flink")
        .withDriverName("org.postgresql.Driver")
        .withUsername("flink")
        .withPassword("secret")
        .build()));

JdbcSink.sink menerima statement SQL, fungsi pemetaan, opsi eksekusi, dan koneksi. Simpan credential di secret manager, bukan hardcode seperti contoh di atas.

Sink Lain: Elasticsearch dan Object Storage

  • Elasticsearch: gunakan ElasticsearchSink untuk pencarian dan dashboarding real-time.
  • Object storage: tulis dengan FileSink ke S3 dalam format Parquet atau ORC untuk data lake.
  • Message queue lain: RabbitMQ dan Kinesis memiliki connector resminya masing-masing di flink-connector-*.

Format Data dan Schema Management

JSON, Avro, Protobuf, CSV

Format data menentukan cara Flink mengubah byte menjadi objek:

Format Avro di Flink SQL
CREATE TABLE orders (
  user_id STRING,
  amount  BIGINT
) WITH (
  'connector' = 'kafka',
  'topic' = 'orders',
  'format' = 'avro'
);

JSON mudah dibaca manusia, CSV ringkas untuk file, Avro dan Protobuf compact serta ber-schema — cocok untuk produksi. Dengan Avro, schema didaftarkan di Schema Registry agar producer dan consumer selalu setuju pada kontrak data.

Memilih Format

  • JSON: cepat untuk debug, overhead besar.
  • Avro: compact, ber-schema, mendukung schema evolution.
  • Protobuf: compact dan strict, populer di tim yang memakai gRPC.
  • CSV: sederhana untuk file dan integrasi legacy.
Memverifikasi konektor yang dimuat
ls $FLINK_HOME/lib | grep -i kafka

Perintah ls $FLINK_HOME/lib | grep -i kafka memastikan JAR konektor sudah berada di folder lib/ — prasyarat agar cluster standalone mengenali connector Kafka.

Penutup

Episode 8 menghubungkan Flink dengan ekosistem: membangun KafkaSource dan KafkaSink, memakai DataGen dan FileSource, menulis ke PostgreSQL dengan JDBC, serta memilih format data JSON, Avro, Protobuf, dan CSV dengan schema management yang rapi.

Inti yang harus dibawa pulang:

  • KafkaSource dan KafkaSink adalah jembatan utama antara Flink dan dunia streaming.
  • Pilih posisi awal pembacaan (earliest atau latest) sesuai kebutuhan job.
  • JdbcSink, ElasticsearchSink, dan FileSink menutupi sebagian besar kebutuhan output.
  • Avro dan Protobuf unggul di produksi karena ber-schema; JSON untuk kecepatan prototipe.
  • Pastikan JAR connector berada di lib/ sebelum menjalankan job di cluster standalone.

Di episode 9 selanjutnya kita akan membahas Table API & SQL — membuat query Flink SQL, memakai TableEnvironment dan catalog, memahami temporal tables dan CDC, serta menerapkan windowing, joins, dan aggregations langsung dari SQL. Ini jalur tercepat menuju produktivitas Flink.

Belajar Apache Flink - Source & Sink Integrations | Belajar Apache Flink