全部產品
Search
文件中心

Simple Log Service:Flink消費

更新時間:Jun 23, 2026

Flink Log Connector 是Log Service提供的 Flink 對接工具,支援開源 Flink 和Realtime Compute Flink 版。本文介紹如何通過 Flink Log Connector 消費和寫入日誌資料。

前提條件

您已完成以下操作:

背景資訊

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)規範實現的新版介面 AliyunLogSourceAliyunLogSink,同時提供 SQL Connector。新作業建議使用新版介面。舊介面 FlinkLogConsumerFlinkLogProducer 後續計劃刪除。

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 方法

是否必填

預設值

含義

setProject(String project)

要消費的Log Service Project。

setLogStore(String logstore)

要消費的 Logstore。

setEndpoint(String endpoint)

Log Service Endpoint,例如 cn-hangzhou.log.aliyuncs.com

setCredentials(String accessKeyId, String accessKey)

訪問Log Service的 AccessKey ID 和 AccessKey Secret。

setDeserializer(AliyunLogDeserializationSchema<T> deserializer)

將 SLS 拉取結果轉換為 Flink 記錄的還原序列化器。

setConsumerGroup(String consumerGroup)

Log Service ConsumerGroup 名稱,用於讀取或提交服務端 checkpoint。

setStartingPosition(StartingPosition)

earliest

消費起始位置。支援 earliestlatestcheckpoint 或 Unix 秒級時間戳記。

setFallbackPosition(StartingPosition)

earliest

起始位置為 checkpoint 且服務端沒有 checkpoint 時使用的兜底位置。

ConfigConstants.LOG_MAX_NUMBER_PER_FETCH

100

單次從單個 Shard 拉取的最大 LogGroup 數量。

ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS

100

一次拉取沒有返回資料時,下次拉取前的等待間隔(單位:毫秒)。

ConfigConstants.LOG_SHARDS_DISCOVERY_INTERVAL_MILLIS

60000

發現 Shard 分裂或合并的輪詢周期(單位:毫秒)。

ConfigConstants.LOG_CHECKPOINT_MODE

ON_CHECKPOINTS

服務端 checkpoint 提交模式。ON_CHECKPOINTS:Flink checkpoint 完成時提交。PERIODIC:獨立定時提交。DISABLED:不提交到服務端。

ConfigConstants.STOP_TIME

停止消費的 Unix 秒級時間戳記,讀取到該時間點後停止消費對應 Shard,適合離線補資料情境。

ConfigConstants.MAX_RETRIES

5

普通錯誤的最大重試次數。

ConfigConstants.SIGNATURE_VERSION

v1

請求籤名版本,支援 v1v4。使用 v4 時需同時設定 REGION_ID

setSplitAssigner(AliyunLogSplitAssigner)

ModuloSplitAssigner

Shard 到 source reader 的分配策略。預設按 shard id 和並行度模數,也可以使用 RoundRobinSplitAssigner 或自訂實現。

ConfigConstants.LOG_COMMIT_INTERVAL_MILLIS

10000

PERIODIC 模式下提交服務端 checkpoint 的間隔(單位:毫秒)。

ConfigConstants.SOURCE_MEMORY_LIMIT

0

拉取結果的記憶體限制(單位:位元組),0 表示不啟用。

ConfigConstants.LOG_USER_AGENT

connector 預設 UA

自訂 User-Agent。

ConfigConstants.REGION_ID

使用 v4 時必填

V4 簽名使用的地區 ID,例如 cn-hangzhou

ConfigConstants.MAX_RETRIES_FOR_RETRYABLE_ERROR

60

可重試錯誤的最大重試次數。

ConfigConstants.BASE_RETRY_BACK_OFF_TIME_MS

200

初始重試退避時間(單位:毫秒)。

ConfigConstants.MAX_RETRY_BACK_OFF_TIME_MS

5000

最大重試退避時間(單位:毫秒)。

ConfigConstants.PROXY_HOST

HTTP Proxy 位址。

ConfigConstants.PROXY_PORT

-1

HTTP 代理連接埠。

ConfigConstants.PROXY_USERNAME

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 方法

是否必填

預設值

含義

setProject(String project)

寫入目標 Project。

setLogStore(String logstore)

預設寫入目標 Logstore。單條 SinkRecord 設定了 logstore 時會覆蓋該值。

setEndpoint(String endpoint)

Log Service Endpoint。

setCredentials(String accessKeyId, String accessKey)

訪問Log Service的 AccessKey ID 和 AccessKey Secret。

setSerializer(AliyunLogSerializationSchema<T> serializer)

將 Flink 記錄轉換為 SinkRecord 的序列化器。

ConfigConstants.FLUSH_INTERVAL_MS

Producer SDK 預設值

日誌在用戶端緩衝後等待發送的最長時間(單位:毫秒)。

ConfigConstants.MAX_RETRIES

Producer SDK 預設值

普通發送失敗的最大重試次數。

ConfigConstants.IO_THREAD_NUM

Producer SDK 預設值

發送日誌的 IO 線程數量。

ConfigConstants.TOTAL_SIZE_IN_BYTES

Producer SDK 預設值

Producer 用戶端可使用的總緩衝大小。

ConfigConstants.MAX_BLOCK_TIME_MS

Producer SDK 預設值

緩衝滿或資源不足時,發送調用最多阻塞等待的時間(單位:毫秒)。

ConfigConstants.SIGNATURE_VERSION

v1

請求籤名版本,支援 v1v4。使用 v4 時需同時設定 REGION_ID

ConfigConstants.BUCKETS

Producer SDK 預設值

Producer 內部分桶數量,用於並發和批量彙總。

ConfigConstants.PRODUCER_ADJUST_SHARD_HASH

true

是否由 Producer 自動調整 shard hash。

ConfigConstants.REGION_ID

使用 v4 時必填

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 參數

適用方向

是否必填

預設值

含義

connector

Source / Sink

固定為 aliyun-log

endpoint

Source / Sink

Log Service Endpoint。

project

Source / Sink

Log Service Project。

logstore

Source / Sink

Source 讀取或 Sink 預設寫入的 Logstore。

access.key.id

Source / Sink

訪問Log Service的 AccessKey ID。

access.key.secret

Source / Sink

訪問Log Service的 AccessKey Secret。

consumer-group

Source

ConsumerGroup 名稱,用於讀取或提交服務端 checkpoint。

scan.startup.mode

Source

earliest

消費起始位置。支援 earliestlatestcheckpoint 或 Unix 秒級時間戳記。

checkpoint.mode

Source

on-checkpoints

服務端 checkpoint 提交模式。支援 on-checkpointsperiodicdisabled

max.number.per.fetch

Source

100

單次從單個 Shard 拉取的最大 LogGroup 數量。

ignore-parse-errors

Source

false

欄位類型轉換失敗時是否輸出 NULL。false 表示拋出異常。

sink.topic

Sink

""

寫入時預設使用的 LogGroup topic,可被 __topic__ 列覆蓋。

sink.source

Sink

寫入時預設使用的 LogGroup source,可被 __source__ 列覆蓋。

flush.interval.ms

Sink

Producer SDK 預設值

日誌在用戶端緩衝後等待發送的最長時間。

signature.version

Source / Sink

v1

請求籤名版本,支援 v1v4

SQL Source 中繼資料列

以下列為內建讀取中繼資料列,聲明後從 SLS log 或 Shard 元資訊中讀取,不從 log content 中取同名欄位:

中繼資料列

推薦類型

含義

__time__

TIMESTAMP(3)

SLS log 時間。

__topic__

STRING

LogGroup topic。

__source__

STRING

LogGroup source。

__shard__

INT

目前記錄所在 Shard ID。

__cursor__

STRING

當前拉取批次對應的 cursor。

SQL Sink 中繼資料列

以下列為內建寫入中繼資料列,聲明後不作為普通 content 寫入:

中繼資料列

推薦類型

含義

__time__

TIMESTAMP(3)

寫入 SLS log time,時間戳記類型按秒寫入。

__topic__

STRING

覆蓋 sink.topic,設定目前記錄的 topic。

__source__

STRING

覆蓋 sink.source,設定目前記錄的 source。

__logstore__

STRING

覆蓋表參數 logstore,將目前記錄寫入指定 Logstore。

__hash_key__

STRING

設定目前記錄的 Shard hash key。

RAM 許可權

使用 Flink Log Connector 訪問Log Service時,需為 RAM 使用者或角色授予以下 API 許可權。

Source 讀取要求的權限

API

Resource

log:GetCursorOrData

acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${logstoreName}

log:ListShards

acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${logstoreName}

log:CreateConsumerGroup

acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${logstoreName}/consumergroup/*

log:ConsumerGroupUpdateCheckPoint

acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${logstoreName}/consumergroup/${consumerGroupName}

Sink 寫入要求的權限

API

Resource

log:PostLogStoreLogs

acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${logstoreName}

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_MILLISConfigConstants.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 中。

  1. 設定啟動參數。

    以下是一個簡單的消費樣本,使用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 產生。

  2. 設定消費起始位置。

    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。

  3. 可選:設定消費進度監控。

    Flink Log Consumer 支援消費進度監控,擷取每個 Shard 的即時消費位置。更多資訊,請參見步驟二:查看消費組狀態

    樣本如下:

    configProps.put(ConfigConstants.LOG_CONSUMERGROUP, "your consumer group name");
    說明

    設定後,Flink Log Consumer 自動建立消費組(已存在則跳過),並將 snapshot 同步到Log Service消費組中。可通過Log Service控制台查看消費進度。

  4. 設定容災和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

  1. 初始化Flink Log Producer。

    1. 初始化 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
    2. 重載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";
            }
        });
  2. 將類比產生的字串寫入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");
    }
}