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.
Crie um projeto e um Logstore do Simple Log Service. Para mais informações, consulte Gerencie projetos e Crie um Logstore básico.
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.
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 |
|
|
Sim |
N/A |
Projeto do Simple Log Service a ser consumido. |
|
|
Sim |
N/A |
Logstore a ser consumido. |
|
|
Sim |
N/A |
Endpoint do Simple Log Service, por exemplo, |
|
|
Sim |
N/A |
AccessKey ID e AccessKey secret usados para acessar o Simple Log Service. |
|
|
Sim |
N/A |
Desserializador que converte os resultados de pull do Simple Log Service em registros do Flink. |
|
|
Não |
N/A |
Nome do grupo de consumidores do Simple Log Service. Usado para ler ou enviar checkpoints no servidor. |
|
|
Não |
|
Posição inicial para consumo. Valores suportados: |
|
|
Não |
|
Posição alternativa usada quando a posição inicial é |
|
|
Não |
|
Número máximo de LogGroups obtidos de um único shard por requisição. |
|
|
Não |
|
Intervalo de espera em milissegundos antes da próxima obtenção quando nenhum dado é retornado. |
|
|
Não |
|
Intervalo de polling em milissegundos para detectar divisões ou mesclagens de shards. |
|
|
Não |
|
Modo de envio de checkpoint no servidor. |
|
|
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. |
|
|
Não |
|
Número máximo de tentativas para erros gerais. |
|
|
Não |
|
Versão da assinatura da requisição. Valores válidos: |
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 |
|
|
Sim |
N/A |
Projeto de destino para gravação. |
|
|
Sim |
N/A |
Logstore de destino padrão para gravação. Este valor é substituído quando um |
|
|
Sim |
N/A |
Endpoint do Simple Log Service. |
|
|
Sim |
N/A |
AccessKey ID e AccessKey secret usados para acessar o Simple Log Service. |
|
|
Sim |
N/A |
Serializador que converte registros do Flink em objetos |
|
|
Não |
Padrão do Producer SDK |
Tempo máximo em milissegundos que os logs ficam em cache no cliente antes do envio. |
|
|
Não |
Padrão do Producer SDK |
Número máximo de tentativas para operações de envio com falha. |
|
|
Não |
Padrão do Producer SDK |
Quantidade de threads de I/O usadas para enviar logs. |
|
|
Não |
Padrão do Producer SDK |
Tamanho total de cache disponível para o cliente Producer. |
|
|
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. |
|
|
Não |
|
Versão da assinatura da requisição. Valores válidos: |
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 |
|
|
Source / Sink |
Sim |
N/A |
Valor fixo: |
|
|
Source / Sink |
Sim |
N/A |
Endpoint do Simple Log Service. |
|
|
Source / Sink |
Sim |
N/A |
Projeto do Simple Log Service. |
|
|
Source / Sink |
Sim |
N/A |
Logstore de leitura (Source) ou Logstore padrão de escrita (Sink). |
|
|
Source / Sink |
Sim |
N/A |
AccessKey ID usado para acessar o Simple Log Service. |
|
|
Source / Sink |
Sim |
N/A |
AccessKey secret usado para acessar o Simple Log Service. |
|
|
Source |
Não |
N/A |
Nome do grupo de consumidores. Usado para ler ou enviar checkpoints no servidor. |
|
|
Source |
Não |
|
Posição inicial para consumo. Valores válidos: |
|
|
Source |
Não |
|
Modo de envio de checkpoint no servidor. Valores válidos: |
|
|
Source |
Não |
|
Número máximo de LogGroups obtidos de um único shard por requisição. |
|
|
Source |
Não |
|
Define se deve retornar |
|
|
Sink |
Não |
|
Tópico padrão do LogGroup usado para gravações. A coluna |
|
|
Sink |
Não |
N/A |
Fonte padrão do LogGroup usada para gravações. A coluna |
|
|
Sink |
Não |
Padrão do Producer SDK |
Tempo máximo que os logs permanecem em cache no cliente antes do envio. |
|
|
Source / Sink |
Não |
|
Versão da assinatura da requisição. Valores válidos: |
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 |
|
|
|
Hora do log do Simple Log Service. |
|
|
|
Tópico do LogGroup. |
|
|
|
Fonte do LogGroup. |
|
|
|
ID do shard do registro atual. |
|
|
|
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 |
|
|
|
Hora do log do Simple Log Service. Um tipo timestamp é gravado em segundos. |
|
|
|
Substitui |
|
|
|
Substitui |
|
|
|
Substitui o parâmetro de tabela |
|
|
|
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 |
|
|
|
|
|
|
|
|
|
|
|
|
Permissões necessárias para gravações no Sink
|
API |
Recurso |
|
|
|
Flink Log Consumer
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_MILLISeConfigConstants.LOG_MAX_NUMBER_PER_FETCHpara 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.
-
Configure os parâmetros de inicialização.
O exemplo a seguir usa
java.util.Propertiespara configuração. Todos os parâmetros são definidos na classeConfigConstants.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. } } } });NotaO 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.
-
Defina a posição inicial para consumo.
Defina
ConfigConstants.LOG_CONSUMER_BEGIN_POSITIONcom 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. UseConfigConstants.LOG_CONSUMERGROUPpara 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);NotaSe o Flink recuperar seu próprio StateBackend, essas configurações serão ignoradas e o consumo será retomado a partir do checkpoint do StateBackend.
-
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");NotaSe 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.
-
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
API legada. Remoção planejada. Use AliyunLogSink para novos jobs.
O Flink Log Producer grava dados no Simple Log Service.
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
-
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_RETRIESSubstitua
LogSerializationSchemae defina um método para serializar dados em umRawLogGroup.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
LogPartitionerpara 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"; } }); -
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");
}
}