Gunakan Flink Log Connector untuk mengonsumsi data log dari Simple Log Service. Connector ini mendukung baik Flink open source maupun Realtime Compute for Apache Flink.
Prasyarat
Simple Log Service telah diaktifkan.
Simple Log Service SDK untuk Python telah diinisialisasi.
-
Proyek dan Logstore Simple Log Service telah dibuat. Kelola proyek. Buat Logstore dasar.
Ikhtisar
Flink Log Connector terdiri dari dua komponen:
-
Konsumen membaca data dari Simple Log Service dengan semantik tepat-sekali (exactly-once semantics) dan load balancing shard.
-
Produsen menulis data ke Simple Log Service.
Tambahkan dependensi Maven berikut ke proyek Anda:
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>flink-log-connector</artifactId>
<version>0.1.38</version>
</dependency>
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>2.5.0</version>
</dependency>
Contoh kode lainnya tersedia di repositori aliyun-log-flink-connector di GitHub.
Flink Log Consumer
Flink Log Consumer berlangganan ke Logstore dan menyediakan semantik tepat-sekali (exactly-once semantics). Consumer ini secara otomatis mendeteksi perubahan shard, sehingga tidak memerlukan pengelolaan shard secara manual.
Setiap subtask mengonsumsi data dari subset shard tertentu. Jika shard dipisah atau digabung, penugasan shard pada subtask akan diperbarui secara otomatis.
Consumer menggunakan operasi API berikut:
-
GetCursorOrData
Mengambil data dari sebuah shard. Pemanggilan yang terlalu sering dapat melebihi batas shard. Gunakan ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS dan ConfigConstants.LOG_MAX_NUMBER_PER_FETCH untuk mengatur interval pemanggilan dan jumlah log yang diambil per pemanggilan. Shard.
Contoh:
configProps.put(ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS, "100"); configProps.put(ConfigConstants.LOG_MAX_NUMBER_PER_FETCH, "100"); -
ListShards
Mengambil semua shard beserta statusnya. Atur interval pemanggilan agar perubahan shard dapat terdeteksi dengan cepat:
// Panggil operasi API ListShards setiap 30 detik. configProps.put(ConfigConstants.LOG_SHARDS_DISCOVERY_INTERVAL_MILLIS, "30000"); -
CreateConsumerGroup
Membuat kelompok konsumen untuk menyinkronkan checkpoint saat Anda mengaktifkan pemantauan progres konsumsi.
-
UpdateCheckPoint
Menyinkronkan snapshot Flink ke kelompok konsumen di Simple Log Service.
-
Konfigurasikan parameter startup.
Contoh berikut menggunakan java.util.Properties untuk konfigurasi. Semua parameter didefinisikan dalam kelas ConfigConstants.
Properties configProps = new Properties(); // Titik akhir Simple Log Service. configProps.put(ConfigConstants.LOG_ENDPOINT, "cn-hangzhou.log.aliyuncs.com"); // Pada contoh ini, ID AccessKey dan Rahasia AccessKey diambil dari variabel lingkungan. String accessKeyId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"); String accessKeySecret = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"); configProps.put(ConfigConstants.LOG_ACCESSKEYID,accessKeyId); configProps.put(ConfigConstants.LOG_ACCESSKEY,accessKeySecret); // Proyek Simple Log Service. String project = "your-project"; // Logstore Simple Log Service. String logstore = "your-logstore"; // Posisi awal untuk mulai mengonsumsi log. configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, Consts.LOG_END_CURSOR); // Metode untuk deserialisasi pesan dari Simple Log Service. FastLogGroupDeserializer deserializer = new FastLogGroupDeserializer(); final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStream<FastLogGroupList> dataStream = env.addSource( new FlinkLogConsumer<FastLogGroupList>(project, logstore, deserializer, configProps) ); dataStream.addSink(new SinkFunction<FastLogGroupList>() { @Override public void invoke(FastLogGroupList logGroupList, Context context) throws Exception { for (FastLogGroup logGroup : logGroupList.getLogGroups()) { int logsCount = logGroup.getLogsCount(); String topic = logGroup.getTopic(); String source = logGroup.getSource(); for (int i = 0; i < logsCount; ++i) { FastLog row = logGroup.getLogs(i); for (int j = 0; j < row.getContentsCount(); ++j) { FastLogContent column = row.getContents(j); // Proses log. System.out.println(column.getKey()); System.out.println(column.getValue()); } } } } }); // Atau, gunakan RawLogGroupListDeserializer. RawLogGroupListDeserializer rawLogGroupListDeserializer = new RawLogGroupListDeserializer(); final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStream<RawLogGroupList> rawLogGroupListDataStream = env.addSource( new FlinkLogConsumer<RawLogGroupList>(project, logstore, rawLogGroupListDeserializer, configProps) ); rawLogGroupListDataStream.addSink(new SinkFunction<RawLogGroupList>() { @Override public void invoke(RawLogGroupList logGroupList, Context context) throws Exception { for (RawLogGroup logGroup : logGroupList.getRawLogGroups()) { String topic = logGroup.getTopic(); String source = logGroup.getSource(); for (RawLog row : logGroup.getLogs()) { // Proses log. } } } });CatatanJumlah subtask tidak bergantung pada jumlah shard. Jika jumlah shard lebih banyak daripada subtask, setiap subtask mengonsumsi data dari satu set shard unik. Jika jumlah shard lebih sedikit, beberapa subtask tetap menganggur hingga shard baru dibuat.
-
Tentukan posisi awal konsumsi.
Atur ConfigConstants.LOG_CONSUMER_BEGIN_POSITION ke salah satu nilai berikut:
-
Consts.LOG_BEGIN_CURSOR: Memulai konsumsi dari awal shard, yaitu data terlama dalam shard tersebut.
-
Consts.LOG_END_CURSOR: Memulai konsumsi dari akhir shard, yaitu data terbaru dalam shard tersebut.
-
Consts.LOG_FROM_CHECKPOINT: Memulai konsumsi dari checkpoint yang disimpan dalam kelompok konsumen tertentu. Gunakan ConfigConstants.LOG_CONSUMERGROUP untuk menentukan kelompok konsumen.
-
UnixTimestamp: String yang merepresentasikan stempel waktu UNIX dalam satuan detik. Konsumsi dimulai dari data yang dicatat setelah stempel waktu ini.
Contoh:
configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, Consts.LOG_BEGIN_CURSOR); configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, Consts.LOG_END_CURSOR); configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, "1512439000"); configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, Consts.LOG_FROM_CHECKPOINT);CatatanJika Flink pulih dari StateBackend-nya sendiri, pengaturan ini diabaikan dan konsumsi dilanjutkan dari checkpoint StateBackend.
-
-
Opsi: Siapkan pemantauan progres konsumsi.
Flink Log Consumer mendukung pemantauan progres konsumsi untuk mengambil posisi konsumsi real-time setiap shard. Langkah 2: Lihat status kelompok konsumen.
Contoh:
configProps.put(ConfigConstants.LOG_CONSUMERGROUP, "your consumer group name");CatatanJika dikonfigurasi, Flink Log Consumer akan membuat kelompok konsumen. Jika kelompok konsumen tersebut sudah ada, tidak ada tindakan yang dilakukan. Snapshot secara otomatis disinkronkan ke kelompok konsumen, dan Anda dapat melihat progres konsumsi di Konsol Simple Log Service.
-
Konfigurasikan pemulihan bencana dan semantik tepat-sekali (exactly-once semantics).
Saat checkpointing Flink diaktifkan, consumer secara berkala menyimpan progres konsumsi. Jika suatu tugas gagal, Flink melanjutkan dari checkpoint terbaru.
Interval checkpoint menentukan jumlah maksimum data yang dikonsumsi ulang saat terjadi kegagalan:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // Aktifkan semantik tepat-sekali Flink. env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // Simpan checkpoint setiap 5 detik. env.enableCheckpointing(5000);Checkpoint dalam dokumentasi Flink.
Flink Log Producer
Flink Log Producer menulis data ke Simple Log Service.
Flink Log Producer hanya mendukung semantik paling sedikit sekali (at-least-once semantics). Data mungkin diduplikasi saat terjadi kegagalan, tetapi tidak ada data yang hilang.
Produsen menggunakan operasi API berikut:
-
PutLogs
-
ListShards
-
Inisialisasi Flink Log Producer.
-
Inisialisasi parameter konfigurasi Properties.
Inisialisasi mirip dengan konsumen. Parameter berikut tersedia (nilai default digunakan jika tidak ditentukan):
// Jumlah thread I/O yang digunakan untuk mengirim data. Nilai default adalah jumlah core CPU. ConfigConstants.IO_THREAD_NUM // Waktu maksimum log dapat di-cache sebelum dikirim. Nilai default adalah 2.000 milidetik. ConfigConstants.FLUSH_INTERVAL_MS // Jumlah total memori yang dapat digunakan oleh suatu tugas. Nilai default adalah 100 MB. ConfigConstants.TOTAL_SIZE_IN_BYTES // Waktu blokir maksimum untuk mengirim log saat batas memori tercapai. Satuan dalam milidetik. Nilai default adalah 60 detik. ConfigConstants.MAX_BLOCK_TIME_MS // Jumlah maksimum percobaan ulang. Nilai default adalah 10. ConfigConstants.MAX_RETRIES -
Timpa LogSerializationSchema dan definisikan metode untuk serialisasi data menjadi RawLogGroup.
RawLogGroup adalah kumpulan log. Detail field tersedia dalam referensi Log.
Untuk menulis data ke shard tertentu, gunakan LogPartitioner untuk menghasilkan kunci hash. Jika tidak dikonfigurasi, data ditulis ke shard acak.
Contohnya:
FlinkLogProducer<String> logProducer = new FlinkLogProducer<String>(new SimpleLogSerializer(), configProps); logProducer.setCustomPartitioner(new LogPartitioner<String>() { // Hasilkan nilai hash 32-bit. public String getHashKey(String element) { try { MessageDigest md = MessageDigest.getInstance("MD5"); md.update(element.getBytes()); String hash = new BigInteger(1, md.digest()).toString(16); while(hash.length() < 32) hash = "0" + hash; return hash; } catch (NoSuchAlgorithmException e) { } return "0000000000000000000000000000000000000000000000000000000000000000"; } });
-
-
Tulis data simulasi ke Simple Log Service:
// Serialisasi data ke format data Simple Log Service. class SimpleLogSerializer implements LogSerializationSchema<String> { public RawLogGroup serialize(String element) { RawLogGroup rlg = new RawLogGroup(); RawLog rl = new RawLog(); rl.setTime((int)(System.currentTimeMillis() / 1000)); rl.addContent("message", element); rlg.addLog(rl); return rlg; } } public class ProducerSample { public static String sEndpoint = "cn-hangzhou.log.aliyuncs.com"; // Pada contoh ini, ID AccessKey dan Rahasia AccessKey diambil dari variabel lingkungan. public static String sAccessKeyId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"); public static String sAccessKey = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"); public static String sProject = "ali-cn-hangzhou-sls-admin"; public static String sLogstore = "test-flink-producer"; private static final Logger LOG = LoggerFactory.getLogger(ConsumerSample.class); public static void main(String[] args) throws Exception { final ParameterTool params = ParameterTool.fromArgs(args); final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.getConfig().setGlobalJobParameters(params); env.setParallelism(3); DataStream<String> simpleStringStream = env.addSource(new EventsGenerator()); Properties configProps = new Properties(); // Titik akhir Simple Log Service. configProps.put(ConfigConstants.LOG_ENDPOINT, sEndpoint); // AccessKey pengguna. configProps.put(ConfigConstants.LOG_ACCESSKEYID, sAccessKeyId); configProps.put(ConfigConstants.LOG_ACCESSKEY, sAccessKey); // Proyek Simple Log Service tempat log ditulis. configProps.put(ConfigConstants.LOG_PROJECT, sProject); // Logstore Simple Log Service tempat log ditulis. configProps.put(ConfigConstants.LOG_LOGSTORE, sLogstore); FlinkLogProducer<String> logProducer = new FlinkLogProducer<String>(new SimpleLogSerializer(), configProps); simpleStringStream.addSink(logProducer); env.execute("flink log producer"); } // Simulasikan pembuatan log. public static class EventsGenerator implements SourceFunction<String> { private boolean running = true; @Override public void run(SourceContext<String> ctx) throws Exception { long seq = 0; while (running) { Thread.sleep(10); ctx.collect((seq++) + "-" + RandomStringUtils.randomAlphabetic(12)); } } @Override public void cancel() { running = false; } } }
Contoh konsumsi
Contoh ini membaca data sebagai FastLogGroupList, mengonversi entri menjadi string JSON dengan flatMap, lalu menulis output ke file teks.
package com.aliyun.openservices.log.flink.sample;
import com.alibaba.fastjson.JSONObject;
import com.aliyun.openservices.log.common.FastLog;
import com.aliyun.openservices.log.common.FastLogGroup;
import com.aliyun.openservices.log.flink.ConfigConstants;
import com.aliyun.openservices.log.flink.FlinkLogConsumer;
import com.aliyun.openservices.log.flink.data.FastLogGroupDeserializer;
import com.aliyun.openservices.log.flink.data.FastLogGroupList;
import com.aliyun.openservices.log.flink.model.CheckpointMode;
import com.aliyun.openservices.log.flink.util.Consts;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.configuration.CheckpointingOptions;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.runtime.state.filesystem.FsStateBackend;
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.CheckpointConfig;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.util.Properties;
public class FlinkConsumerSample {
private static final String SLS_ENDPOINT = "your-endpoint";
// Pada contoh ini, ID AccessKey dan Rahasia AccessKey diambil dari variabel lingkungan.
private static final String ACCESS_KEY_ID = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
private static final String ACCESS_KEY_SECRET = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
private static final String SLS_PROJECT = "your-project";
private static final String SLS_LOGSTORE = "your-logstore";
public static void main(String[] args) throws Exception {
final ParameterTool params = ParameterTool.fromArgs(args);
Configuration conf = new Configuration();
// Direktori checkpoint seperti "file:///tmp/flink"
conf.setString(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "your-checkpoint-dir");
final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(1, conf);
env.getConfig().setGlobalJobParameters(params);
env.setParallelism(1);
env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
env.setStateBackend(new FsStateBackend("file:///tmp/flinkstate"));
Properties configProps = new Properties();
configProps.put(ConfigConstants.LOG_ENDPOINT, SLS_ENDPOINT);
configProps.put(ConfigConstants.LOG_ACCESSKEYID, ACCESS_KEY_ID);
configProps.put(ConfigConstants.LOG_ACCESSKEY, ACCESS_KEY_SECRET);
configProps.put(ConfigConstants.LOG_MAX_NUMBER_PER_FETCH, "10");
configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, Consts.LOG_FROM_CHECKPOINT);
configProps.put(ConfigConstants.LOG_CONSUMERGROUP, "your-consumer-group");
configProps.put(ConfigConstants.LOG_CHECKPOINT_MODE, CheckpointMode.ON_CHECKPOINTS.name());
configProps.put(ConfigConstants.LOG_COMMIT_INTERVAL_MILLIS, "10000");
FastLogGroupDeserializer deserializer = new FastLogGroupDeserializer();
DataStream<FastLogGroupList> stream = env.addSource(
new FlinkLogConsumer<>(SLS_PROJECT, SLS_LOGSTORE, deserializer, configProps));
stream.flatMap((FlatMapFunction<FastLogGroupList, String>) (value, out) -> {
for (FastLogGroup logGroup : value.getLogGroups()) {
int logCount = logGroup.getLogsCount();
for (int i = 0; i < logCount; i++) {
FastLog log = logGroup.getLogs(i);
JSONObject jsonObject = new JSONObject();
jsonObject.put("topic", logGroup.getTopic());
jsonObject.put("source", logGroup.getSource());
for (int j = 0; j < log.getContentsCount(); j++) {
jsonObject.put(log.getContents(j).getKey(), log.getContents(j).getValue());
}
out.collect(jsonObject.toJSONString());
}
}
}).returns(String.class);
stream.writeAsText("log-" + System.nanoTime());
env.execute("Flink consumer");
}
}