Todos os produtos
Search
Central de documentação

Simple Log Service:Consumir dados usando o Flink

Última atualização: Aug 26, 2026

Use o Flink Log Connector para consumir dados de log do Simple Log Service. O conector oferece suporte tanto ao Flink open source quanto ao Realtime Compute for Apache Flink.

Pré-requisitos

  • O Simple Log Service está ativado.

  • O SDK do Simple Log Service para Python está inicializado.

Visão geral

O Flink Log Connector possui dois componentes:

  • O consumidor lê dados do Simple Log Service com semântica exactly-once e balanceamento de carga entre shards.

  • O produtor grava dados no Simple Log Service.

Adicione as seguintes dependências Maven ao seu projeto:

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

Mais exemplos de código estão disponíveis no repositório aliyun-log-flink-connector no GitHub.

Nota

A versão 0.1.46 do Flink Log Connector introduz AliyunLogSource e AliyunLogSink, novas interfaces baseadas na especificação FLIP-27, além de um SQL Connector. Use as novas interfaces para novos jobs. As interfaces legadas FlinkLogConsumer e FlinkLogProducer têm remoção planejada.

AliyunLogSource (DataStream Source)

AliyunLogSource baseia-se na especificação FLIP-27 e integra-se ao Flink por meio de env.fromSource(...). Os estados de split e cursor da Source participam dos checkpoints do Flink para recuperação de falhas em jobs. Se um ConsumerGroup estiver configurado, os checkpoints também podem ser enviados ao servidor do Simple Log Service para monitoramento do progresso de consumo.

Uso básico:

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");

Parâmetros da Source

Parâmetro / Método do Builder

Obrigatório

Padrão

Descrição

setProject(String project)

Sim

N/A

Projeto do Simple Log Service a ser consumido.

setLogStore(String logstore)

Sim

N/A

Logstore a ser consumido.

setEndpoint(String endpoint)

Sim

N/A

Endpoint do Simple Log Service, por exemplo, cn-hangzhou.log.aliyuncs.com.

setCredentials(String accessKeyId, String accessKey)

Sim

N/A

AccessKey ID e AccessKey secret usados para acessar o Simple Log Service.

setDeserializer(AliyunLogDeserializationSchema<T> deserializer)

Sim

N/A

Desserializador que converte os resultados de pull do Simple Log Service em registros do Flink.

setConsumerGroup(String consumerGroup)

Não

N/A

Nome do grupo de consumidores do Simple Log Service. Usado para ler ou enviar checkpoints no servidor.

setStartingPosition(StartingPosition)

Não

earliest

Posição inicial para consumo. Valores suportados: earliest, latest, checkpoint ou timestamp Unix em segundos.

setFallbackPosition(StartingPosition)

Não

earliest

Posição alternativa usada quando a posição inicial é checkpoint e não existe checkpoint no servidor.

ConfigConstants.LOG_MAX_NUMBER_PER_FETCH

Não

100

Número máximo de LogGroups obtidos de um único shard por requisição.

ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS

Não

100

Intervalo de espera em milissegundos antes da próxima obtenção quando nenhum dado é retornado.

ConfigConstants.LOG_SHARDS_DISCOVERY_INTERVAL_MILLIS

Não

60000

Intervalo de polling em milissegundos para detectar divisões ou mesclagens de shards.

ConfigConstants.LOG_CHECKPOINT_MODE

Não

ON_CHECKPOINTS

Modo de envio de checkpoint no servidor. ON_CHECKPOINTS: envia quando um checkpoint do Flink é concluído. PERIODIC: envia em um intervalo separado. DISABLED: não envia ao servidor.

ConfigConstants.STOP_TIME

Não

N/A

Timestamp Unix em segundos. O consumo para o shard correspondente após este momento. Este parâmetro é útil para preenchimento retroativo de dados offline.

ConfigConstants.MAX_RETRIES

Não

5

Número máximo de tentativas para erros gerais.

ConfigConstants.SIGNATURE_VERSION

Não

v1

Versão da assinatura da requisição. Valores válidos: v1 e v4. Se você usar v4, também deve definir REGION_ID.

Desserializador personalizado

Implemente a interface AliyunLogDeserializationSchema<T> e conclua a expansão de logs e conversão de campos no método deserialize. Um PullLogsResult pode conter vários LogGroups, e cada LogGroup pode conter múltiplos logs. O desserializador pode gerar zero, um ou mais registros do Flink para o Collector.

O exemplo a seguir expande cada log do Simple Log Service em um POJO contendo metadados e um mapa de conteúdo:

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));
            }
        }
    }
}

Exemplo completo de consumo

Este exemplo obtém credenciais de acesso de variáveis de ambiente e retoma o consumo a partir de um checkpoint de ConsumerGroup no servidor. Caso não exista checkpoint no servidor, o consumo inicia na posição mais antiga.

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, "30000");
        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 integra-se ao Flink por meio de stream.sinkTo(...). O Sink utiliza o Producer SDK do Simple Log Service para enviar dados de forma assíncrona e aguarda a conclusão das requisições enviadas durante os checkpoints do Flink ou encerramento do job, fornecendo semântica at-least-once.

O serializador personalizado deve implementar a interface AliyunLogSerializationSchema<T>. Um único elemento de entrada pode gerar zero, um ou mais registros do Simple Log Service via Collector<SinkRecord>.

Uso básico:

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");

Parâmetros do Sink

Parâmetro / Método do Builder

Obrigatório

Padrão

Descrição

setProject(String project)

Sim

N/A

Projeto de destino para gravação.

setLogStore(String logstore)

Sim

N/A

Logstore de destino padrão para gravação. Este valor é substituído quando um SinkRecord especifica um Logstore diferente.

setEndpoint(String endpoint)

Sim

N/A

Endpoint do Simple Log Service.

setCredentials(String accessKeyId, String accessKey)

Sim

N/A

AccessKey ID e AccessKey secret usados para acessar o Simple Log Service.

setSerializer(AliyunLogSerializationSchema<T> serializer)

Sim

N/A

Serializador que converte registros do Flink em objetos SinkRecord.

ConfigConstants.FLUSH_INTERVAL_MS

Não

Padrão do Producer SDK

Tempo máximo em milissegundos que os logs ficam em cache no cliente antes do envio.

ConfigConstants.MAX_RETRIES

Não

Padrão do Producer SDK

Número máximo de tentativas para operações de envio com falha.

ConfigConstants.IO_THREAD_NUM

Não

Padrão do Producer SDK

Quantidade de threads de I/O usadas para enviar logs.

ConfigConstants.TOTAL_SIZE_IN_BYTES

Não

Padrão do Producer SDK

Tamanho total de cache disponível para o cliente Producer.

ConfigConstants.MAX_BLOCK_TIME_MS

Não

Padrão do Producer SDK

Tempo máximo em milissegundos que uma chamada de envio bloqueia quando o cache está cheio ou os recursos são insuficientes.

ConfigConstants.SIGNATURE_VERSION

Não

v1

Versão da assinatura da requisição. Valores válidos: v1 e v4. Se você usar v4, também deve definir REGION_ID.

Para controlar em qual shard os dados são gravados, chame record.setHashKey(...) no SinkRecord para definir uma chave hash. Para gravar dinamicamente em diferentes Logstores, chame record.setLogstore(...) para substituir o Logstore padrão do Sink.

Exemplo completo de gravação

Este exemplo usa env.fromSequence(...) para gerar dados de teste e gravá-los no Simple Log Service por meio de AliyunLogSerializationSchema.

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

O identificador do SQL Connector é aliyun-log. O mesmo conector oferece suporte tanto a SQL Source quanto a SQL Sink. Em uma tabela Source, colunas regulares são lidas do conteúdo de log do Simple Log Service pela correspondência de nomes de colunas. Em uma tabela Sink, colunas regulares são gravadas no conteúdo do log pelos nomes das colunas.

Exemplo de 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'
);

Exemplo de 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'
);

Parâmetros SQL WITH

Parâmetro SQL

Direção aplicável

Obrigatório

Padrão

Descrição

connector

Source / Sink

Sim

N/A

Valor fixo: aliyun-log.

endpoint

Source / Sink

Sim

N/A

Endpoint do Simple Log Service.

project

Source / Sink

Sim

N/A

Projeto do Simple Log Service.

logstore

Source / Sink

Sim

N/A

Logstore de leitura (Source) ou Logstore padrão de escrita (Sink).

access.key.id

Source / Sink

Sim

N/A

AccessKey ID usado para acessar o Simple Log Service.

access.key.secret

Source / Sink

Sim

N/A

AccessKey secret usado para acessar o Simple Log Service.

consumer-group

Source

Não

N/A

Nome do grupo de consumidores. Usado para ler ou enviar checkpoints no servidor.

scan.startup.mode

Source

Não

earliest

Posição inicial para consumo. Valores válidos: earliest, latest, checkpoint ou timestamp Unix em segundos.

checkpoint.mode

Source

Não

on-checkpoints

Modo de envio de checkpoint no servidor. Valores válidos: on-checkpoints, periodic e disabled.

max.number.per.fetch

Source

Não

100

Número máximo de LogGroups obtidos de um único shard por requisição.

ignore-parse-errors

Source

Não

false

Define se deve retornar NULL quando a conversão de tipo de campo falhar. O valor false lança uma exceção.

sink.topic

Sink

Não

""

Tópico padrão do LogGroup usado para gravações. A coluna __topic__ pode substituir este valor.

sink.source

Sink

Não

N/A

Fonte padrão do LogGroup usada para gravações. A coluna __source__ pode substituir este valor.

flush.interval.ms

Sink

Não

Padrão do Producer SDK

Tempo máximo que os logs permanecem em cache no cliente antes do envio.

signature.version

Source / Sink

Não

v1

Versão da assinatura da requisição. Valores válidos: v1 e v4.

Colunas de metadados do SQL Source

Os seguintes nomes de colunas são metadados integrados para leitura. Quando declarados, os valores são lidos dos metadados de log ou shard do Simple Log Service, em vez de um campo de conteúdo de log com o mesmo nome.

Coluna de metadados

Tipo recomendado

Descrição

__time__

TIMESTAMP(3)

Hora do log do Simple Log Service.

__topic__

STRING

Tópico do LogGroup.

__source__

STRING

Fonte do LogGroup.

__shard__

INT

ID do shard do registro atual.

__cursor__

STRING

Cursor do lote de pull atual.

Colunas de metadados do SQL Sink

Os nomes de colunas a seguir são metadados integrados para gravação. Quando declaradas, essas colunas não são escritas como conteúdo regular.

Coluna de metadados

Tipo recomendado

Descrição

__time__

TIMESTAMP(3)

Hora do log do Simple Log Service. Um tipo timestamp é gravado em segundos.

__topic__

STRING

Substitui sink.topic para definir o tópico do registro atual.

__source__

STRING

Substitui sink.source para definir a fonte do registro atual.

__logstore__

STRING

Substitui o parâmetro de tabela logstore para gravar o registro atual no Logstore especificado.

__hash_key__

STRING

Define a chave hash do shard para o registro atual.

Permissões RAM

Ao usar o Flink Log Connector para acessar o Simple Log Service, conceda as permissões de API correspondentes ao usuário ou função RAM.

Permissões necessárias para leituras da Source

API

Recurso

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}

Permissões necessárias para gravações no Sink

API

Recurso

log:PostLogStoreLogs

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

Flink Log Consumer

Importante

API legada. Remoção planejada. Use AliyunLogSource para novos jobs.

O Flink Log Consumer assina um Logstore e fornece semântica exactly-once. Ele detecta alterações de shards automaticamente, eliminando o gerenciamento manual de shards.

Cada subtarefa consome dados de um subconjunto de shards. Se os shards forem divididos ou mesclados, a atribuição de shards da subtarefa será atualizada automaticamente.

O consumidor utiliza as seguintes operações de API:

  • GetCursorOrData

    Recupera dados de um shard. Chamadas frequentes podem exceder o limite do shard. Use ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS e ConfigConstants.LOG_MAX_NUMBER_PER_FETCH para controlar o intervalo de chamadas e a quantidade de logs recuperados por chamada. Para mais informações, consulte Shards.

    Exemplo:

    configProps.put(ConfigConstants.LOG_FETCH_DATA_INTERVAL_MILLIS, "100");
    configProps.put(ConfigConstants.LOG_MAX_NUMBER_PER_FETCH, "100");
  • ListShards

    Recupera todos os shards e seus status. Ajuste o intervalo de chamadas para detectar alterações de shards prontamente:

    // Call the ListShards API operation every 30s.
    configProps.put(ConfigConstants.LOG_SHARDS_DISCOVERY_INTERVAL_MILLIS, "30000");
  • CreateConsumerGroup

    Cria um grupo de consumidores para sincronizar checkpoints quando você ativa o monitoramento de progresso de consumo.

  • UpdateCheckPoint

    Sincroniza snapshots do Flink com o grupo de consumidores no Simple Log Service.

  1. Configure os parâmetros de inicialização.

    O exemplo a seguir usa java.util.Properties para configuração. Todos os parâmetros são definidos na classe ConfigConstants.

    Properties configProps = new Properties();
    // The endpoint of Simple Log Service.
    configProps.put(ConfigConstants.LOG_ENDPOINT, "cn-hangzhou.log.aliyuncs.com");
    // In this example, the AccessKey ID and AccessKey secret are obtained from environment variables.
    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);
    // The Simple Log Service project.
    String project = "your-project";
    // The Simple Log Service Logstore.
    String logstore = "your-logstore";
    // The position from which to start consuming logs.
    configProps.put(ConfigConstants.LOG_CONSUMER_BEGIN_POSITION, Consts.LOG_END_CURSOR);
    // The method to deserialize messages from 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);
                        // Process logs.
                        System.out.println(column.getKey());
                        System.out.println(column.getValue());
                    }
                }
            }
        }
    });
    // Or, use 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()) {
                    // Process logs.
                }
            }
        }
    });
    Nota

    O número de subtarefas é independente da contagem de shards. Se houver mais shards do que subtarefas, cada subtarefa consome dados de um conjunto exclusivo de shards. Se houver menos shards, algumas subtarefas permanecerão ociosas até que novos shards sejam criados.

  2. Defina a posição inicial para consumo.

    Defina ConfigConstants.LOG_CONSUMER_BEGIN_POSITION com um dos seguintes valores:

    • Consts.LOG_BEGIN_CURSOR: Inicia o consumo do início do shard, que corresponde aos dados mais antigos no shard.

    • Consts.LOG_END_CURSOR: Inicia o consumo do final do shard, que corresponde aos dados mais recentes no shard.

    • Consts.LOG_FROM_CHECKPOINT: Inicia o consumo a partir de um checkpoint salvo em um grupo de consumidores específico. Use ConfigConstants.LOG_CONSUMERGROUP para especificar o grupo de consumidores.

    • UnixTimestamp: Uma string que representa um timestamp UNIX em segundos. O consumo começa a partir dos dados registrados após este timestamp.

    Exemplo:

    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);
    Nota

    Se o Flink recuperar seu próprio StateBackend, essas configurações serão ignoradas e o consumo será retomado a partir do checkpoint do StateBackend.

  3. Opcional: Configure o monitoramento de progresso de consumo.

    O Flink Log Consumer oferece suporte ao monitoramento de progresso de consumo para recuperar a posição de consumo em tempo real de cada shard. Para mais informações, consulte Etapa 2: Visualize o status de um grupo de consumidores.

    Exemplo:

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

    Se configurado, o Flink Log Consumer cria um grupo de consumidores. Se o grupo já existir, nenhuma ação será tomada. Snapshots são sincronizados automaticamente com o grupo de consumidores, e você pode visualizar o progresso de consumo no console do Simple Log Service.

  4. Configure recuperação de desastres e semântica exactly-once.

    Quando o checkpointing do Flink está ativado, o consumidor salva periodicamente o progresso de consumo. Se uma tarefa falhar, o Flink retoma a partir do último checkpoint.

    O intervalo de checkpoint determina quantos dados são reprocessados em caso de falha:

    final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    // Enable Flink exactly-once semantics.
    env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
    // Save a checkpoint every 5s.
    env.enableCheckpointing(5000);

    Para mais informações sobre checkpoints do Flink, consulte Checkpoints na documentação do Flink.

Flink Log Producer

Importante

API legada. Remoção planejada. Use AliyunLogSink para novos jobs.

O Flink Log Producer grava dados no Simple Log Service.

Nota

O Flink Log Producer oferece suporte apenas à semântica at-least-once. Dados podem ser duplicados em caso de falha, mas nenhum dado é perdido.

O produtor utiliza as seguintes operações de API:

  • PutLogs

  • ListShards

  1. Inicialize o Flink Log Producer.

    Inicialize os parâmetros de configuração Properties.

    A inicialização é semelhante à do consumidor. Os seguintes parâmetros estão disponíveis (padrões são usados se não especificados):

    // The number of I/O threads used to send data. The default value is the number of CPU cores.
    ConfigConstants.IO_THREAD_NUM
    // The maximum time that logs can be cached before being sent. The default value is 2,000 milliseconds.
    ConfigConstants.FLUSH_INTERVAL_MS
    // The total amount of memory that a task can use. The default value is 100 MB.
    ConfigConstants.TOTAL_SIZE_IN_BYTES
    // The maximum blocking time for sending logs when the memory limit is reached. The unit is milliseconds. The default value is 60s.
    ConfigConstants.MAX_BLOCK_TIME_MS
    // The maximum number of retries. The default value is 10.
    ConfigConstants.MAX_RETRIES

    Substitua LogSerializationSchema e defina um método para serializar dados em um RawLogGroup.

    Um RawLogGroup é uma coleção de logs. Para mais informações sobre os campos, consulte Log.

    Para gravar dados em um shard específico, use LogPartitioner para gerar uma chave hash. Se você não configurar um particionador, os dados serão gravados em shards aleatórios.

    Por exemplo:

    FlinkLogProducer<String> logProducer = new FlinkLogProducer<String>(new SimpleLogSerializer(), configProps);
    logProducer.setCustomPartitioner(new LogPartitioner<String>() {
          // Generate a 32-bit hash value.
          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. Grave dados simulados no Simple Log Service:

    // Serialize data into the Simple Log Service data format.
    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";
        // In this example, the AccessKey ID and AccessKey secret are obtained from environment variables.
        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();
            // The endpoint of Simple Log Service.
            configProps.put(ConfigConstants.LOG_ENDPOINT, sEndpoint);
            // The user's AccessKey.
            configProps.put(ConfigConstants.LOG_ACCESSKEYID, sAccessKeyId);
            configProps.put(ConfigConstants.LOG_ACCESSKEY, sAccessKey);
            // The Simple Log Service project to which logs are written.
            configProps.put(ConfigConstants.LOG_PROJECT, sProject);
            // The Simple Log Service Logstore to which logs are written.
            configProps.put(ConfigConstants.LOG_LOGSTORE, sLogstore);
            FlinkLogProducer<String> logProducer = new FlinkLogProducer<String>(new SimpleLogSerializer(), configProps);
            simpleStringStream.addSink(logProducer);
            env.execute("flink log producer");
        }
        // Simulate log generation.
        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;
            }
        }
    }

Exemplo de consumo

Este exemplo lê dados como FastLogGroupList, converte entradas em strings JSON com flatMap e grava a saída em um arquivo de texto.

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";
    // In this example, the AccessKey ID and AccessKey secret are obtained from environment variables.
    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");
    }
}