Flink Log Connector 是Log Service提供的 Flink 對接工具,支援開源 Flink 和Realtime Compute Flink 版。本文介紹如何通過 Flink Log Connector 消費和寫入日誌資料。
前提條件
您已完成以下操作:
已建立 Project 和 Logstore。具體操作,請參見管理Project和建立基礎LogStore。
背景資訊
Flink Log Connector 包含消費者(Consumer)和生產者(Producer)兩部分:
消費者(Consumer):從Log Service讀取資料,支援 exactly once 語義和 Shard 負載平衡。
生產者(Producer):將資料寫入Log Service。
使用 Flink Log Connector 時,需在專案中添加以下 Maven 依賴:
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>flink-log-connector</artifactId>
<version>0.1.46</version>
</dependency>
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>2.5.0</version>
</dependency>更多程式碼範例,請參見aliyun-log-flink-connector。
Flink Log Connector 0.1.46 版本引入了基於 FLIP-27(Flink 改進提案 27)規範實現的新版介面 AliyunLogSource 和 AliyunLogSink,同時提供 SQL Connector。新作業建議使用新版介面。舊介面 FlinkLogConsumer 和 FlinkLogProducer 後續計劃刪除。
AliyunLogSource(DataStream Source)
AliyunLogSource 基於 FLIP-27 規範實現,通過 env.fromSource(...) 接入。Source 的 split/cursor 狀態參與 Flink checkpoint,用於作業 failover 恢複。如果設定了 ConsumerGroup,還可以將 checkpoint 提交到Log Service服務端,便於監控消費進度。
基本用法:
Properties properties = new Properties();
properties.setProperty(ConfigConstants.LOG_CHECKPOINT_MODE, CheckpointMode.ON_CHECKPOINTS.name());
properties.setProperty(ConfigConstants.LOG_MAX_NUMBER_PER_FETCH, "100");
AliyunLogSource<MyRecord> source = AliyunLogSource.<MyRecord>builder()
.setProject("your-project")
.setLogStore("your-logstore")
.setEndpoint("cn-hangzhou.log.aliyuncs.com")
.setCredentials(accessKeyId, accessKeySecret)
.setConsumerGroup("flink-source-consumer")
.setStartingPosition(StartingPosition.EARLIEST)
.setProperties(properties)
.setDeserializer(new MyDeserializer())
.build();
DataStream<MyRecord> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"aliyun-log-source");Source 參數
參數 / Builder 方法 | 是否必填 | 預設值 | 含義 |
| 是 | 無 | 要消費的Log Service Project。 |
| 是 | 無 | 要消費的 Logstore。 |
| 是 | 無 | Log Service Endpoint,例如 |
| 是 | 無 | 訪問Log Service的 AccessKey ID 和 AccessKey Secret。 |
| 是 | 無 | 將 SLS 拉取結果轉換為 Flink 記錄的還原序列化器。 |
| 否 | 無 | Log Service ConsumerGroup 名稱,用於讀取或提交服務端 checkpoint。 |
| 否 |
| 消費起始位置。支援 |
| 否 |
| 起始位置為 checkpoint 且服務端沒有 checkpoint 時使用的兜底位置。 |
| 否 |
| 單次從單個 Shard 拉取的最大 LogGroup 數量。 |
| 否 |
| 一次拉取沒有返回資料時,下次拉取前的等待間隔(單位:毫秒)。 |
| 否 |
| 發現 Shard 分裂或合并的輪詢周期(單位:毫秒)。 |
| 否 |
| 服務端 checkpoint 提交模式。 |
| 否 | 無 | 停止消費的 Unix 秒級時間戳記,讀取到該時間點後停止消費對應 Shard,適合離線補資料情境。 |
| 否 |
| 普通錯誤的最大重試次數。 |
| 否 |
| 請求籤名版本,支援 |
| 否 |
| Shard 到 source reader 的分配策略。預設按 shard id 和並行度模數,也可以使用 |
| 否 |
|
|
| 否 |
| 拉取結果的記憶體限制(單位:位元組), |
| 否 | connector 預設 UA | 自訂 User-Agent。 |
| 使用 | 無 | V4 簽名使用的地區 ID,例如 |
| 否 |
| 可重試錯誤的最大重試次數。 |
| 否 |
| 初始重試退避時間(單位:毫秒)。 |
| 否 |
| 最大重試退避時間(單位:毫秒)。 |
| 否 | 無 | HTTP Proxy 位址。 |
| 否 |
| HTTP 代理連接埠。 |
| 否 | 無 | HTTP 代理使用者名稱。 |
自訂 Deserializer
實現 AliyunLogDeserializationSchema<T> 介面,在 deserialize 方法中完成日誌展開和欄位轉換。一個 PullLogsResult 可能包含多個 LogGroup,每個 LogGroup 可能包含多條日誌,還原序列化器可以向 Collector 輸出零條、一條或多條 Flink 記錄。
以下樣本將每條 SLS 日誌展開為包含中繼資料和 content map 的 POJO:
public class ContentMapDeserializer implements AliyunLogDeserializationSchema<SlsLogRecord> {
@Override
public TypeInformation<SlsLogRecord> getProducedType() {
return TypeInformation.of(SlsLogRecord.class);
}
@Override
public void deserialize(PullLogsResult record, Collector<SlsLogRecord> out) {
for (LogGroupData logGroupData : record.getLogGroupList()) {
FastLogGroup logGroup = logGroupData.GetFastLogGroup();
for (int logIndex = 0; logIndex < logGroup.getLogsCount(); logIndex++) {
FastLog log = logGroup.getLogs(logIndex);
Map<String, String> fields = new LinkedHashMap<>();
for (int contentIndex = 0; contentIndex < log.getContentsCount(); contentIndex++) {
FastLogContent content = log.getContents(contentIndex);
fields.put(content.getKey(), content.getValue());
}
out.collect(new SlsLogRecord(
log.getTime(),
logGroup.getTopic(),
logGroup.getSource(),
record.getShard(),
record.getCursor(),
fields));
}
}
}
}完整消費樣本
以下樣本從環境變數擷取訪問憑證,從 ConsumerGroup checkpoint 繼續消費。服務端無 checkpoint 時,從最早位置開始。
package com.aliyun.openservices.log.flink.sample;
import com.aliyun.openservices.log.flink.ConfigConstants;
import com.aliyun.openservices.log.flink.model.CheckpointMode;
import com.aliyun.openservices.log.flink.source.AliyunLogSource;
import com.aliyun.openservices.log.flink.source.StartingPosition;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.configuration.CheckpointingOptions;
import org.apache.flink.configuration.Configuration;
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 AliyunLogConsumerSample {
private static final String SLS_ENDPOINT = "cn-hangzhou.log.aliyuncs.com";
private static final String SLS_PROJECT = "your-project";
private static final String SLS_LOGSTORE = "your-logstore";
private static final String CONSUMER_GROUP = "your-consumer-group";
public static void main(String[] args) throws Exception {
String accessKeyId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
String accessKeySecret = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
Configuration configuration = new Configuration();
configuration.setString(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "file:///tmp/flink-checkpoints");
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(configuration);
env.setParallelism(2);
env.enableCheckpointing(60000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().enableExternalizedCheckpoints(
CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
Properties sourceProperties = new Properties();
sourceProperties.setProperty(ConfigConstants.LOG_MAX_NUMBER_PER_FETCH, "100");
sourceProperties.setProperty(ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS, "100");
sourceProperties.setProperty(ConfigConstants.LOG_SHARDS_DISCOVERY_INTERVAL_MILLIS, "60000");
sourceProperties.setProperty(ConfigConstants.LOG_CHECKPOINT_MODE, CheckpointMode.ON_CHECKPOINTS.name());
AliyunLogSource<SlsLogRecord> source = AliyunLogSource.<SlsLogRecord>builder()
.setEndpoint(SLS_ENDPOINT)
.setProject(SLS_PROJECT)
.setLogStore(SLS_LOGSTORE)
.setCredentials(accessKeyId, accessKeySecret)
.setConsumerGroup(CONSUMER_GROUP)
.setStartingPosition(StartingPosition.CHECKPOINT)
.setFallbackPosition(StartingPosition.EARLIEST)
.setProperties(sourceProperties)
.setDeserializer(new ContentMapDeserializer())
.build();
DataStream<SlsLogRecord> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"aliyun-log-source");
stream.print();
env.execute("aliyun log consumer");
}
}AliyunLogSink(DataStream Sink)
AliyunLogSink 通過 stream.sinkTo(...) 接入。Sink 基於Log Service Producer SDK 非同步發送資料,並在 Flink checkpoint 或作業結束時等待已提交請求完成,提供 at-least-once 語義。
自訂序列化器需實現 AliyunLogSerializationSchema<T> 介面。每個輸入元素可通過 Collector<SinkRecord> 輸出零條、一條或多條 SLS 記錄。
基本用法:
class MySerializationSchema implements AliyunLogSerializationSchema<String> {
@Override
public void serialize(String element, Collector<SinkRecord> output) {
LogItem item = new LogItem((int) (System.currentTimeMillis() / 1000L));
item.PushBack("message", element);
SinkRecord record = new SinkRecord();
record.setTopic("flink");
record.setSource("flink-job");
record.setLogItem(item);
output.collect(record);
}
}
AliyunLogSink<String> sink = AliyunLogSink.<String>builder()
.setProject("your-project")
.setLogStore("your-logstore")
.setEndpoint("cn-hangzhou.log.aliyuncs.com")
.setCredentials(accessKeyId, accessKeySecret)
.setSerializer(new MySerializationSchema())
.setProperty(ConfigConstants.FLUSH_INTERVAL_MS, "100")
.build();
stream.sinkTo(sink).name("aliyun-log-sink");Sink 參數
參數 / Builder 方法 | 是否必填 | 預設值 | 含義 |
| 是 | 無 | 寫入目標 Project。 |
| 是 | 無 | 預設寫入目標 Logstore。單條 |
| 是 | 無 | Log Service Endpoint。 |
| 是 | 無 | 訪問Log Service的 AccessKey ID 和 AccessKey Secret。 |
| 是 | 無 | 將 Flink 記錄轉換為 |
| 否 | Producer SDK 預設值 | 日誌在用戶端緩衝後等待發送的最長時間(單位:毫秒)。 |
| 否 | Producer SDK 預設值 | 普通發送失敗的最大重試次數。 |
| 否 | Producer SDK 預設值 | 發送日誌的 IO 線程數量。 |
| 否 | Producer SDK 預設值 | Producer 用戶端可使用的總緩衝大小。 |
| 否 | Producer SDK 預設值 | 緩衝滿或資源不足時,發送調用最多阻塞等待的時間(單位:毫秒)。 |
| 否 |
| 請求籤名版本,支援 |
| 否 | Producer SDK 預設值 | Producer 內部分桶數量,用於並發和批量彙總。 |
| 否 |
| 是否由 Producer 自動調整 shard hash。 |
| 使用 | 無 | V4 簽名使用的地區 ID。 |
如果需要控制寫入 Shard,可以在 SinkRecord 上調用 record.setHashKey(...) 設定 hash key。如果需要動態寫入不同 Logstore,可以調用 record.setLogstore(...) 覆蓋 Sink 預設 Logstore。
完整寫入樣本
以下樣本使用 env.fromSequence(...) 產生測試資料,並通過 AliyunLogSerializationSchema 寫入Log Service。
package com.aliyun.openservices.log.flink.sample;
import com.aliyun.openservices.log.common.LogItem;
import com.aliyun.openservices.log.flink.ConfigConstants;
import com.aliyun.openservices.log.flink.data.SinkRecord;
import com.aliyun.openservices.log.flink.model.AliyunLogSerializationSchema;
import com.aliyun.openservices.log.flink.sink.AliyunLogSink;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;
public class AliyunLogProducerSample {
private static final String SLS_ENDPOINT = "cn-hangzhou.log.aliyuncs.com";
private static final String SLS_PROJECT = "your-project";
private static final String SLS_LOGSTORE = "your-logstore";
public static void main(String[] args) throws Exception {
String accessKeyId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
String accessKeySecret = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(3);
DataStream<Long> events = env.fromSequence(1, 1000);
AliyunLogSink<Long> sink = AliyunLogSink.<Long>builder()
.setEndpoint(SLS_ENDPOINT)
.setProject(SLS_PROJECT)
.setLogStore(SLS_LOGSTORE)
.setCredentials(accessKeyId, accessKeySecret)
.setSerializer(new LongSerializer())
.setProperty(ConfigConstants.FLUSH_INTERVAL_MS, "100")
.setProperty(ConfigConstants.MAX_RETRIES, "10")
.build();
events.sinkTo(sink).name("aliyun-log-sink");
env.execute("aliyun log producer");
}
public static class LongSerializer implements AliyunLogSerializationSchema<Long> {
@Override
public void serialize(Long element, Collector<SinkRecord> output) {
LogItem logItem = new LogItem((int) (System.currentTimeMillis() / 1000L));
logItem.PushBack("id", String.valueOf(element));
logItem.PushBack("message", "message-" + element);
SinkRecord record = new SinkRecord();
record.setTopic("flink");
record.setSource("flink-job");
record.setLogItem(logItem);
output.collect(record);
}
}
}SQL Connector
SQL Connector 標識為 aliyun-log,同一 connector 同時支援 SQL Source 和 SQL Sink。Source 普通列按同名 SLS log content 讀取,Sink 普通列按列名寫入 SLS log content。
SQL Source 樣本
CREATE TABLE sls_logs (
`__time__` TIMESTAMP(3),
`__topic__` STRING,
`__source__` STRING,
level STRING,
message STRING,
status_code INT
) WITH (
'connector' = 'aliyun-log',
'endpoint' = 'cn-hangzhou.log.aliyuncs.com',
'project' = 'your-project',
'logstore' = 'your-logstore',
'access.key.id' = '${ACCESS_KEY_ID}',
'access.key.secret' = '${ACCESS_KEY_SECRET}',
'consumer-group' = 'flink-sql-consumer',
'scan.startup.mode' = 'checkpoint',
'scan.startup.default-position' = 'earliest',
'checkpoint.mode' = 'on-checkpoints',
'max.number.per.fetch' = '100',
'shards.discovery.interval.ms' = '60000',
'ignore-parse-errors' = 'true'
);SQL Sink 樣本
CREATE TABLE sls_sink (
`__time__` TIMESTAMP(3),
`__topic__` STRING,
`__source__` STRING,
level STRING,
message STRING,
status_code INT
) WITH (
'connector' = 'aliyun-log',
'endpoint' = 'cn-hangzhou.log.aliyuncs.com',
'project' = 'your-project',
'logstore' = 'your-logstore',
'access.key.id' = '${ACCESS_KEY_ID}',
'access.key.secret' = '${ACCESS_KEY_SECRET}',
'sink.topic' = 'flink-sql',
'sink.source' = 'flink-job',
'flush.interval.ms' = '100',
'max.retries' = '5'
);SQL WITH 參數
SQL 參數 | 適用方向 | 是否必填 | 預設值 | 含義 |
| Source / Sink | 是 | 無 | 固定為 |
| Source / Sink | 是 | 無 | Log Service Endpoint。 |
| Source / Sink | 是 | 無 | Log Service Project。 |
| Source / Sink | 是 | 無 | Source 讀取或 Sink 預設寫入的 Logstore。 |
| Source / Sink | 是 | 無 | 訪問Log Service的 AccessKey ID。 |
| Source / Sink | 是 | 無 | 訪問Log Service的 AccessKey Secret。 |
| Source | 否 | 無 | ConsumerGroup 名稱,用於讀取或提交服務端 checkpoint。 |
| Source | 否 |
| 消費起始位置。支援 |
| Source | 否 |
| 服務端 checkpoint 提交模式。支援 |
| Source | 否 |
| 單次從單個 Shard 拉取的最大 LogGroup 數量。 |
| Source | 否 |
| 欄位類型轉換失敗時是否輸出 NULL。 |
| Sink | 否 |
| 寫入時預設使用的 LogGroup topic,可被 |
| Sink | 否 | 無 | 寫入時預設使用的 LogGroup source,可被 |
| Sink | 否 | Producer SDK 預設值 | 日誌在用戶端緩衝後等待發送的最長時間。 |
| Source / Sink | 否 |
| 請求籤名版本,支援 |
SQL Source 中繼資料列
以下列為內建讀取中繼資料列,聲明後從 SLS log 或 Shard 元資訊中讀取,不從 log content 中取同名欄位:
中繼資料列 | 推薦類型 | 含義 |
|
| SLS log 時間。 |
|
| LogGroup topic。 |
|
| LogGroup source。 |
|
| 目前記錄所在 Shard ID。 |
|
| 當前拉取批次對應的 cursor。 |
SQL Sink 中繼資料列
以下列為內建寫入中繼資料列,聲明後不作為普通 content 寫入:
中繼資料列 | 推薦類型 | 含義 |
|
| 寫入 SLS log time,時間戳記類型按秒寫入。 |
|
| 覆蓋 |
|
| 覆蓋 |
|
| 覆蓋表參數 |
|
| 設定目前記錄的 Shard hash key。 |
RAM 許可權
使用 Flink Log Connector 訪問Log Service時,需為 RAM 使用者或角色授予以下 API 許可權。
Source 讀取要求的權限
API | Resource |
|
|
|
|
|
|
|
|
Sink 寫入要求的權限
API | Resource |
|
|
Flink Log Consumer
該介面後續計劃刪除。新作業請使用 AliyunLogSource。
Flink Log Consumer 用於訂閱 Logstore 中的日誌資料,支援 exactly once 語義。Flink Log Consumer 自動感知 Shard 分裂和合并,無需手動處理 Shard 變化。
Flink 的每個子任務負責消費 Logstore 中的部分 Shard。當 Shard 發生分裂或合并時,子任務消費的 Shard 會自動調整。
Flink Log Consumer 使用的Log Service API 介面:
GetCursorOrData
用於從Shard中擷取資料,注意頻繁的調用該介面可能會導致資料超過Log Service的Shard限額,可以通過ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS和ConfigConstants.LOG_MAX_NUMBER_PER_FETCH控制介面調用的時間間隔和每次調用擷取的日誌數量。Shard的限額請參見分區(Shard)。
樣本如下:
configProps.put(ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS, "100"); configProps.put(ConfigConstants.LOG_MAX_NUMBER_PER_FETCH, "100");ListShards
擷取 Logstore 中所有 Shard 列表及狀態。Shard 經常分裂合并時,可縮短輪詢周期及時發現變化。樣本:
// 設定每60s調用一次ListShards介面。 configProps.put(ConfigConstants.LOG_SHARDS_DISCOVERY_INTERVAL_MILLIS, "60000");CreateConsumerGroup
設定消費進度監控時調用,建立 ConsumerGroup 用於同步 Checkpoint。
UpdateCheckPoint
將 Flink 的 snapshot 同步到Log Service ConsumerGroup 中。
設定啟動參數。
以下是一個簡單的消費樣本,使用java.util.Properties作為組態工具,所有Flink Log Consumer的配置均在ConfigConstants中。
Properties configProps = new Properties(); // 設定訪問Log Service的網域名稱。 configProps.put(ConfigConstants.LOG_ENDPOINT, "cn-hangzhou.log.aliyuncs.com"); // 本樣本從環境變數中擷取AccessKey ID和AccessKey Secret。 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); // 設定Log Service的project。 String project = "your-project"; // 設定Log Service的LogStore。 String logstore = "your-logstore"; // 設定消費Log Service起始位置。 configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, Consts.LOG_END_CURSOR); // 設定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); // 處理日誌 System.out.println(column.getKey()); System.out.println(column.getValue()); } } } } }); // 或者使用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()) { // 處理日誌 } } } });說明Flink 子任務數量和 Logstore 的 Shard 數量相互獨立。Shard 多於子任務時,每個子任務不重複地消費 Shard;Shard 少於子任務時,部分子任務空閑,直到新 Shard 產生。
設定消費起始位置。
Flink Log Consumer支援設定Shard的消費起始位置,通過設定屬性ConfigConstants.LOG_CONSUMER_BEGIN_POSITION,就可以定製從Shard的頭、尾或者某個特定時間開始消費。另外,Flink Log Connector也支援從某個具體的消費組中恢複消費。屬性的具體取值如下:
Consts.LOG_BEGIN_CURSOR:從 Shard 頭部(最舊資料)開始消費。
Consts.LOG_END_CURSOR:從 Shard 尾部(最新資料)開始消費。
Consts.LOG_FROM_CHECKPOINT:表示從某個特定的消費組中儲存的Checkpoint開始消費,通過ConfigConstants.LOG_CONSUMERGROUP指定具體的消費組。
UnixTimestamp:Unix 秒級時間戳記(整型數值字串),從該時間點之後開始消費。
樣本如下:
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);說明從 Flink StateBackend 恢複時,Flink Log Connector 忽略上述設定,使用 StateBackend 中儲存的 Checkpoint。
可選:設定消費進度監控。
Flink Log Consumer 支援消費進度監控,擷取每個 Shard 的即時消費位置。更多資訊,請參見步驟二:查看消費組狀態。
樣本如下:
configProps.put(ConfigConstants.LOG_CONSUMERGROUP, "your consumer group name");說明設定後,Flink Log Consumer 自動建立消費組(已存在則跳過),並將 snapshot 同步到Log Service消費組中。可通過Log Service控制台查看消費進度。
設定容災和exactly once語義支援。
開啟 Flink Checkpointing 後,Flink Log Consumer 周期性儲存每個 Shard 的消費進度。任務失敗時,從最新 Checkpoint 恢複消費。
Checkpoint 周期決定任務失敗時最多回溯的資料量。配置樣本:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 開啟Flink exactly once語義。 env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // 每5s儲存一次Checkpoint。 env.enableCheckpointing(5000);
Flink Log Producer
該介面後續計劃刪除。新作業請使用 AliyunLogSink。
Flink Log Producer 將資料寫入Log Service。
Flink Log Producer 僅支援 at least once 語義,任務失敗時寫入的資料可能重複但不會丟失。
Flink Log Producer 使用的Log Service API 介面:
PutLogs
ListShards
初始化Flink Log Producer。
初始化 Properties 設定參數。
Flink Log Producer 初始化方式與 Consumer 類似,以下參數使用預設值即可,按需自訂:
// 用於發送資料的I/O線程的數量,預設為核心數。 ConfigConstants.IO_THREAD_NUM // 日誌發送前被緩衝的最大允許時間,預設為2000毫秒。 ConfigConstants.FLUSH_INTERVAL_MS // 任務可以使用的記憶體總的大小,預設為100 MB。 ConfigConstants.TOTAL_SIZE_IN_BYTES // 記憶體達到上限時,發送日誌的最大阻塞時間,單位為毫秒,預設為 60s。 ConfigConstants.MAX_BLOCK_TIME_MS // 最大重試次數,預設為 10 次。 ConfigConstants.MAX_RETRIES重載LogSerializationSchema,定義將資料序列化成RawLogGroup的方法。
RawLogGroup 是日誌集合,欄位含義請參見日誌(Log)。
如需將資料寫入指定 Shard,使用 LogPartitioner 產生 HashKey。未設定 LogPartitioner 時,資料隨機寫入 Shard。
樣本如下:
FlinkLogProducer<String> logProducer = new FlinkLogProducer<String>(new SimpleLogSerializer(), configProps); logProducer.setCustomPartitioner(new LogPartitioner<String>() { // 產生32位Hash值。 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"; } });
將類比產生的字串寫入Log Service,樣本如下:
// 將資料序列化成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"; //本樣本從環境變數中擷取AccessKey ID和AccessKey Secret。 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(); // 設定訪問Log Service的網域名稱。 configProps.put(ConfigConstants.LOG_ENDPOINT, sEndpoint); // 設定使用者AK。 configProps.put(ConfigConstants.LOG_ACCESSKEYID, sAccessKeyId); configProps.put(ConfigConstants.LOG_ACCESSKEY, sAccessKey); // 設定日誌寫入的Log Serviceproject。 configProps.put(ConfigConstants.LOG_PROJECT, sProject); // 設定日誌寫入的Log ServiceLogstore。 configProps.put(ConfigConstants.LOG_LOGSTORE, sLogstore); FlinkLogProducer<String> logProducer = new FlinkLogProducer<String>(new SimpleLogSerializer(), configProps); simpleStringStream.addSink(logProducer); env.execute("flink log producer"); } // 類比產生日誌。 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; } } }
消費樣本
以下樣本使用 Flink Log Consumer 讀取資料,通過 flatMap 將 FastLogGroupList 轉換為 JSON 字串,輸出到命令列或寫入文字檔。
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";
//本樣本從環境變數中擷取AccessKey ID和AccessKey Secret。
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();
// Checkpoint dir like "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");
}
}