Saiba como usar o conector do Log Service (SLS).
Contexto
O Simple Log Service é um serviço de ponta a ponta para dados de log. Ele auxilia na coleta, no consumo, no envio, na consulta e na análise eficiente desses dados, melhorando a eficiência de O&M e permitindo processar grandes volumes de informações.
A tabela a seguir lista as capacidades do conector SLS.
|
Categoria |
Descrição |
|
Tipos suportados |
Tabela de origem e tabela de destino |
|
Modo de execução |
Apenas modo streaming |
|
Métricas específicas do conector |
N/A |
|
Formato de dados |
N/A |
|
Tipo de API |
SQL, DataStream API e API YAML de ingestão de dados |
|
Atualização ou exclusão de dados na tabela de destino |
As tabelas de destino são apenas de adição (append-only); não é possível atualizar ou excluir dados. |
Recursos
O conector de source SLS lê diretamente os campos de atributos das mensagens. A tabela abaixo detalha os campos suportados.
|
Parâmetro |
Tipo |
Descrição |
|
__source__ |
STRING METADATA VIRTUAL |
Origem da mensagem. |
|
__topic__ |
STRING METADATA VIRTUAL |
Tópico da mensagem. |
|
__timestamp__ |
BIGINT METADATA VIRTUAL |
Horário do log. |
|
__tag__ |
MAP<VARCHAR, VARCHAR> METADATA VIRTUAL |
Tag da mensagem. Por exemplo, para o atributo |
Pré-requisitos
Certifique-se de ter criado um projeto do Log Service e um Logstore. Para mais informações, consulte Create a Project and a Logstore.
Limitações
Somente o Ververica Runtime (VVR) 11,1 e versões posteriores suportam o uso do SLS como fonte síncrona para ingestão de dados definida em yaml.
O conector SLS garante apenas a semântica at-least-once.
Não defina o
source parallelismsuperior ao número de shards, pois isso desperdiça recursos. Além disso, no Ververica Runtime (VVR) 8.0.5 e anteriores, uma alteração na quantidade de shards pode causar falha no recurso automático defailover, impedindo o consumo de alguns shards.
SQL
Sintaxe
CREATE TABLE sls_table(
a INT,
b INT,
c VARCHAR
) WITH (
'connector' = 'sls',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'logStore' = '<yourLogStoreName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
Opções With
-
Gerais
Parâmetro
Descrição
Tipo
Obrigatório
Padrão
Observações
connector
Conector a ser usado.
String
Sim
Nenhum
Defina como
sls.endPoint
Endpoint do Log Service (SLS).
String
Sim
Nenhum
Especifique o endereço de acesso VPC do Log Service (SLS). Para mais informações, consulte Service endpoints.
Nota-
Por padrão, o Realtime Compute for Apache Flink não acessa a internet. Para habilitar o acesso à internet a partir da sua Virtual Private Cloud (VPC), use um NAT Gateway. Para mais detalhes, veja How do I access the Internet?.
-
Recomendamos não acessar o SLS pela internet. Caso seja necessário, use HTTPS e ative a aceleração de transferência. Consulte Manage transfer acceleration para saber mais.
project
Nome do projeto SLS.
String
Sim
Nenhum
N/A
logStore
Nome do Logstore ou MetricStore.
String
Sim
Nenhum
Os dados em um Logstore são consumidos da mesma forma que os dados em um MetricStore.
accessId
AccessKey ID da sua conta Alibaba Cloud.
String
Sim
Nenhum
Para mais informações, consulte Obtain an AccessKey pair.
ImportantePara evitar a exposição do seu par AccessKey, recomendamos usar variáveis para especificar o AccessKey ID e o AccessKey secret. Veja Project variables para detalhes.
accessKey
AccessKey secret da sua conta Alibaba Cloud.
String
Sim
Nenhum
-
-
Específicas da source
Parâmetro
Descrição
Tipo
Obrigatório
Padrão
Observações
enableNewSource
Define se deve ser usada a nova fonte de dados que implementa a interface FLIP-27.
Boolean
Não
false
A nova fonte adapta-se automaticamente a mudanças nos shards e distribui-os da forma mais equilibrada possível entre todas as subtarefas de source.
Importante-
Esta opção é suportada apenas no VVR 8.0.9 e posteriores. O valor padrão é
truepara VVR 11,1 e superiores. -
Se você alterar o valor desta opção, não será possível restaurar jobs a partir de um estado salvo. Para retomar o consumo de um offset histórico, inicie primeiro o job com a opção
consumerGrouppara registrar o progresso de consumo no grupo de consumidores do SLS. Em seguida, defina a opçãoconsumeFromCheckpointcomotruee reinicie o job sem estado. -
Caso um Logstore possua shards somente leitura, algumas subtarefas podem continuar solicitando dados de outros shards após concluírem os seus próprios. Isso pode gerar desequilíbrio na carga de trabalho e afetar o desempenho. Para mitigar esse problema, ajuste o paralelismo, otimize a estratégia de agendamento ou mescle shards pequenos para reduzir a quantidade total e simplificar a alocação de tarefas.
shardDiscoveryIntervalMs
Intervalo para descoberta dinâmica de shards.
Long
Não
60000
Defina esta opção com um valor negativo para desativar a descoberta dinâmica. Unidade: milissegundos.
Nota-
O valor deve ser maior ou igual a 60.000 milissegundos (1 minuto).
-
Esta opção só tem efeito quando
enableNewSourceestá definido comotrue. -
Suportado apenas no VVR 8.0.9 e versões posteriores.
startupMode
Modo de inicialização da tabela de source.
String
Não
timestamp
-
timestamp(padrão): Consome logs a partir do horário inicial especificado. -
latest: Inicia o consumo de logs pelo offset mais recente. -
earliest: Começa o consumo de logs pelo offset mais antigo. -
consumer_group: Inicia o consumo a partir do offset registrado pelo grupo de consumidores. Se o grupo não tiver um offset registrado para determinado shard, o consumo começa pelo offset mais antigo.
Importante-
Para versões do VVR anteriores à 11,1, o valor consumer_group não é suportado. É necessário definir
consumeFromCheckpointcomotrue. Nesse caso, o consumo de logs inicia no offset registrado pelo grupo de consumidores especificado, e a configuração do modo de inicialização não terá efeito.
startTime
Horário inicial para consumo de logs.
String
Não
Horário atual
O formato é
yyyy-MM-dd hh:mm:ss.Tem efeito apenas quando
startupModeestá definido comotimestamp.NotaAs opções
startTimeestopTimebaseiam-se no atributo__receive_time__do SLS, e não no atributo__timestamp__.stopTime
Horário final para consumo de logs.
String
Não
Nenhum
O formato é
yyyy-MM-dd hh:mm:ss.Nota-
Esta opção serve apenas para consumir logs históricos e deve ser configurada com uma data passada. Se definida para um horário futuro, o consumo poderá parar prematuramente caso nenhum novo log seja gravado, causando interrupção no fluxo de dados sem mensagens de erro.
-
Para que o job Flink encerre após o consumo de todos os logs, também é necessário definir
exitAfterFinishcomotrue.
consumerGroup
Nome do grupo de consumidores.
String
Não
Nenhum
Um grupo de consumidores registra o progresso do consumo. Você pode definir um nome personalizado sem formato fixo.
NotaJobs Flink diferentes devem usar grupos de consumidores distintos. Se múltiplos jobs utilizarem o mesmo grupo, eles não se coordenarão e cada um consumirá todos os dados. Isso ocorre porque o Flink não usa o grupo de consumidores do SLS para atribuição de partições ao consumir dados. Consequentemente, cada consumidor processa mensagens independentemente, mesmo compartilhando o mesmo grupo.
consumeFromCheckpoint
Define se o consumo deve iniciar a partir de um checkpoint do grupo de consumidores.
String
Não
false
-
true: Também é necessário especificar um grupo de consumidores. O programa Flink inicia o consumo dos logs a partir do checkpoint salvo nesse grupo. Caso não exista um checkpoint correspondente, o consumo começa pelo valor configurado em startTime. -
false(valor padrão): Não inicia o consumo de logs a partir do checkpoint salvo para o grupo de consumidores especificado.
ImportanteEste parâmetro não é mais suportado no VVR 11,1 e versões posteriores. Para essas versões, defina a opção
startupModecomoconsumer_group.maxRetries
Número de tentativas após falha na leitura do SLS.
String
Não
3
N/A
batchGetSize
Quantidade de grupos de logs lidos por requisição.
String
Não
100
A configuração
batchGetSizenão pode exceder 1000, caso contrário um erro será reportado.exitAfterFinish
Define se o job Flink encerra após todo o dado ser consumido.
String
Não
false
-
true: O programa Flink encerra após o consumo completo dos dados. -
false(padrão): O programa Flink não encerra após a conclusão do consumo de dados.
query
ImportanteEsta opção foi descontinuada no VVR 11.3, mas versões posteriores permanecem compatíveis.
Instrução de consulta para pré-processamento de dados antes do consumo.
String
Não
Nenhum
Use esta opção para filtrar dados do SLS antes que o Flink os consuma. Isso reduz custos e aumenta a velocidade de processamento.
Por exemplo,
'query' = '*| where request_method = ''GET'''indica que, antes de ler dados do SLS, o Flink selecionará apenas registros onde o camporequest_methodtenha o valor 'GET'.NotaEsta opção utiliza a linguagem SPL do Log Service (SLS). Para mais informações, consulte SPL syntax.
Importante-
Suportado apenas no VVR 8.0.1 e versões posteriores.
-
Este recurso gera cobranças do Log Service (SLS). Consulte Billing of Log Service para detalhes.
processor
Nome do processador SLS para pré-processamento de dados. Se tanto
queryquantoprocessorforem especificados,queryterá precedência eprocessorserá ignorado.String
Não
Nenhum
Esta opção filtra dados do SLS antes do consumo pelo Flink, reduzindo custos e melhorando a velocidade de processamento. Recomendamos o uso de
processorem vez dequery.Por exemplo,
'processor' = 'test-filter-processor'indica que o processador SLS filtrará os dados antes que o Flink faça a leitura.NotaEsta opção utiliza a linguagem SPL do Log Service (SLS). Para mais informações, consulte SPL syntax. Para saber como criar ou atualizar um processador SLS, veja Manage processors.
ImportanteSuportado apenas no VVR 11.3 e versões posteriores.
Este recurso gera cobranças do Log Service (SLS). Consulte Billing of Log Service para detalhes.
preserveRawBytes
Define se os bytes brutos transportados pelo SLS devem ser buscados e preservados diretamente quando
enableNewSourceestiver definido comotrue.Boolean
Não
false
Quando esta opção está ativada, campos BINARY/VARBINARY leem o byte[] bruto diretamente, e campos CHAR/VARCHAR constroem strings a partir dos bytes brutos. Outros campos continuam usando a lógica de conversão padrão. O comportamento é consistente com o da fonte antiga.
Nota-
Suportado apenas no VVR 11.8 e versões posteriores.
-
Esta opção só tem efeito quando
enableNewSourceestá definido comotrue.
-
-
Específicas do destino
Parâmetro
Descrição
Tipo
Obrigatório
Padrão
Observações
topicField
Especifica um campo cujo valor sobrescreve o atributo
__topic__, indicando o tópico do log.String
Não
Nenhum
O valor desta opção deve ser um campo existente na tabela.
timeField
Especifica um campo cujo valor sobrescreve o atributo
__timestamp__, indicando o horário de gravação do log.String
Não
Horário atual
O valor desta opção deve ser um campo
INTexistente na tabela. Se não especificado, o horário atual será usado.sourceField
Especifica um campo cujo valor sobrescreve o atributo
__source__, indicando a source do log, como o endereço IP da máquina que o gerou.String
Não
Nenhum
O valor desta opção deve ser um campo existente na tabela.
partitionField
Define um campo para particionamento. Um hash do valor deste campo determina qual shard receberá os dados, garantindo que registros com o mesmo hash vão para o mesmo shard.
String
Não
Nenhum
Se esta opção não for especificada, cada registro será gravado aleatoriamente em um shard disponível.
buckets
Se
partitionFieldfor especificado, esta opção define o número de buckets para mapeamento de valores hash.String
Não
64
O valor deve ser uma potência de 2 no intervalo [1, 256]. O número de buckets deve ser maior ou igual ao número de shards. Caso contrário, alguns shards podem não receber dados.
flushIntervalMs
Intervalo de gravação de dados.
String
Não
2000
Unidade: milissegundos.
writeNullProperties
Define se valores nulos devem ser gravados como strings vazias no SLS.
Boolean
Não
true
-
true(valor padrão): Grava valores nulos no log como strings vazias. -
false: Campos avaliados como nulo não são gravados no log.
NotaSuportado apenas no VVR 8.0.6 e versões posteriores.
-
Mapeamento de tipos
|
Tipo Flink |
Tipo SLS |
|
BOOLEAN |
STRING |
|
VARBINARY |
|
|
VARCHAR |
|
|
TINYINT |
|
|
INTEGER |
|
|
BIGINT |
|
|
FLOAT |
|
|
DOUBLE |
|
|
DECIMAL |
Ingestão de dados (Beta)
Limitações
Este recurso é suportado apenas pelo Realtime Compute for Apache Flink versão 11,1 e posteriores.
Sintaxe
source:
type: sls
name: SLS Source
endpoint: <endpoint>
project: <project>
logstore: <logstore>
accessId: <accessId>
accessKey: <accessKey>
Parâmetros
|
Parâmetro |
Descrição |
Tipo |
Obrigatório |
Padrão |
Observações |
||||||||||
|
type |
Tipo da fonte de dados. |
String |
Sim |
Nenhum |
O valor deve ser |
||||||||||
|
endpoint |
Endpoint do Log Service (SLS). |
String |
Sim |
Nenhum |
Endereço de acesso VPC do Log Service (SLS). Para mais informações, consulte Service Endpoints. Nota
|
||||||||||
|
accessId |
AccessKey ID da sua conta Alibaba Cloud. |
String |
Sim |
Nenhum |
Para mais informações, consulte How do I view the AccessKey ID and AccessKey secret information?. Importante
Para evitar a exposição das suas informações AccessKey, recomendamos usar uma variável de projeto para especificar o valor do AccessKey. Veja Project variables para mais detalhes. |
||||||||||
|
accessKey |
AccessKey secret da sua conta Alibaba Cloud. |
String |
Sim |
Nenhum |
|||||||||||
|
project |
Nome do projeto Log Service (SLS). |
String |
Sim |
Nenhum |
Nenhum |
||||||||||
|
logStore |
Nome do Logstore ou Metricstore do SLS. |
String |
Sim |
Nenhum |
Os dados em um Logstore são consumidos da mesma forma que os dados em um Metricstore. |
||||||||||
|
schema.inference.strategy |
Estratégia de inferência de schema. |
String |
Não |
continuous |
|
||||||||||
|
maxPreFetchLogGroups |
Número máximo de grupos de logs a serem lidos de cada shard para a inferência inicial de schema. |
Integer |
Não |
50 |
Antes de ler e processar dados, o conector pré-consome uma quantidade especificada de grupos de logs de cada shard para inicializar as informações de schema. |
||||||||||
|
shardDiscoveryIntervalMs |
Intervalo, em milissegundos, para descoberta dinâmica de alterações nos shards. |
Long |
Não |
60000 |
Defina este parâmetro com um valor negativo para desativar a descoberta dinâmica. Nota
O valor deve ser maior ou igual a 60.000 milissegundos (1 minuto). |
||||||||||
|
startupMode |
Modo de inicialização. |
String |
Não |
Nenhum |
|
||||||||||
|
startTime |
Horário inicial para consumo de logs. |
String |
Não |
Horário atual |
O formato é Este parâmetro só tem efeito quando Nota
Os parâmetros |
||||||||||
|
stopTime |
Horário final para consumo de logs. |
String |
Não |
Nenhum |
O formato é Nota
Para que o job Flink encerre após o consumo de todos os logs, também é necessário definir |
||||||||||
|
consumerGroup |
Nome do grupo de consumidores. |
String |
Não |
Nenhum |
Um grupo de consumidores registra o progresso do consumo. Você pode especificar qualquer nome personalizado. |
||||||||||
|
batchGetSize |
Número de grupos de logs lidos por requisição. |
Integer |
Não |
100 |
O valor de batchGetSize não pode exceder 1.000. Caso contrário, ocorrerá um erro. |
||||||||||
|
maxRetries |
Número de tentativas em caso de falha na leitura do Log Service (SLS). |
Integer |
Não |
3 |
Nenhum |
||||||||||
|
exitAfterFinish |
Define se o job Flink encerra após o consumo de todos os dados. |
Boolean |
Não |
false |
|
||||||||||
|
query |
Instrução de pré-processamento para consumo de dados do Log Service (SLS). |
String |
Não |
Nenhum |
Use este parâmetro para filtrar dados no Log Service (SLS) antes do consumo, economizando custos e aumentando a velocidade de processamento. Por exemplo, Nota
A consulta deve usar a sintaxe SPL do Log Service. Para mais informações, consulte SPL syntax. Importante
|
||||||||||
|
compressType |
Tipo de compressão para o Log Service (SLS). |
String |
Não |
Nenhum |
Tipos de compressão suportados incluem:
|
||||||||||
|
timeZone |
Fuso horário para |
String |
Não |
Nenhum |
Por padrão, nenhum deslocamento é adicionado. |
||||||||||
|
regionId |
Região onde o Log Service (SLS) está implantado. |
String |
Não |
Nenhum |
Para mais informações, consulte Supported regions. |
||||||||||
|
signVersion |
Versão da assinatura da requisição para o Log Service (SLS). |
String |
Não |
Nenhum |
Para mais informações, consulte Request signatures. |
||||||||||
|
shardModDivisor |
Divisor usado na leitura de shards do Logstore. |
Int |
Não |
-1 |
Para mais informações, consulte Shards. |
||||||||||
|
shardModRemainder |
Resto usado na leitura de shards do Logstore. |
Int |
Não |
-1 |
Para mais informações, consulte Shards. |
||||||||||
|
metadata.list |
Colunas de metadados a serem passadas para jobs downstream. |
String |
Não |
Nenhum |
Os campos de metadados disponíveis incluem |
||||||||||
|
decode.table-id.fields |
Especifica campos cujos valores são usados para gerar um Table ID ao analisar dados de log do Log Service (SLS). |
String |
Não |
Nenhum |
Múltiplos campos são separados por vírgula
Nota
Este parâmetro é suportado no Realtime Compute for Apache Flink versão 11,6 e posteriores. |
||||||||||
|
fixed-types |
Especifica os tipos de dados para campos específicos ao analisar logs do Log Service (SLS). |
String |
Não |
Nenhum |
Ao analisar dados, especifique os tipos para campos específicos. Use vírgula Nota
Este parâmetro é suportado no Realtime Compute for Apache Flink versão 11,6 e posteriores. |
||||||||||
|
timestamp-format.standard |
Formato para campos de timestamp nos dados de log do Log Service (SLS). |
String |
Não |
SQL |
Valores válidos:
Nota
Este parâmetro é suportado no Realtime Compute for Apache Flink versão 11,6 e posteriores. |
||||||||||
|
ingestion.ignore-errors |
Define se erros ocorridos durante a análise de dados devem ser ignorados. |
Boolean |
Não |
false |
Nota
Este parâmetro é suportado no Realtime Compute for Apache Flink versão 11,6 e posteriores. |
||||||||||
|
ingestion.error-tolerance.max-count |
Se |
Integer |
Não |
-1 |
Este parâmetro só tem efeito quando Nota
Este parâmetro é suportado no Realtime Compute for Apache Flink versão 11,6 e posteriores. |
Reutilizar um catálogo existente
A partir da versão 11,5 do Realtime Compute for Apache Flink, é possível referenciar um Catálogo SLS integrado criado na página Data Management em jobs de ingestão de dados Flink CDC. Isso elimina a necessidade de especificar manualmente as propriedades de conexão.
source:
type: sls
using.built-in-catalog: sls_catalog
Atualmente, os jobs de ingestão de dados podem reutilizar automaticamente os seguintes parâmetros do Catálogo SLS integrado:
endpoint
project
accessId
accessKey
Para substituir esses parâmetros reutilizados automaticamente, defina-os explicitamente na configuração yaml. Os parâmetros definidos no arquivo yaml têm precedência.
Mapeamento de tipos de dados
Se fixed-types não estiver configurado, aplica-se o seguinte mapeamento de tipos de dados:
|
Tipo SLS |
Tipo CDC |
|
STRING |
STRING |
Se fixed-types estiver configurado, o sistema analisa os dados utilizando os tipos especificados.
Inferência e evolução de schema
-
Pré-consumo de dados de shard e inicialização de schema
O conector SLS mantém o schema do Logstore que está lendo. Antes de ler os dados, o conector pré-consome até
maxPreFetchLogGroupsgrupos de logs de cada shard. Ele analisa o schema de cada entrada de log e os mescla para inicializar o schema da tabela. Um evento de criação de tabela é então gerado com base nesse schema inicial antes do início do consumo de dados.NotaPara cada shard, o conector tenta consumir dados começando de uma hora antes do horário atual para analisar o schema do log.
-
Informações de chave primária
Logs do Log Service (SLS) não contêm informações de chave primária. Adicione manualmente uma chave primária à tabela usando regras de transformação:
transform: - source-table: <project>.<logstore> projection: * primary-keys: key1, key2 -
Inferência de schema e alterações de schema
Após a inicialização do schema, se schema.inference.strategy estiver definido como
static, o conector SLS analisa cada entrada de log com base no schema inicial e não gera eventos de alteração de schema. Se schema.inference.strategy estiver definido comocontinuous, o conector analisa cada entrada, infere as colunas físicas e as compara com o schema atual. Se o schema inferido for inconsistente com o atual, os schemas são mesclados conforme as seguintes regras:Se as colunas físicas inferidas contiverem campos que não estão no schema atual, o conector adiciona esses campos ao schema e gera um evento para adicionar colunas anuláveis.
Se as colunas físicas inferidas não tiverem campos que existem no schema atual, o conector retém esses campos, preenche seus dados com NULL e não gera evento de exclusão de coluna.
O conector SLS infere o tipo de dados de todos os campos em cada entrada de log como String. Atualmente, apenas a adição de novas colunas é suportada. O conector anexa novas colunas ao final do schema e as define como anuláveis.
Exemplos de código
-
SQL para tabelas de source e destino
CREATE TEMPORARY TABLE sls_input( `time` BIGINT, url STRING, dt STRING, float_field FLOAT, double_field DOUBLE, boolean_field BOOLEAN, `__topic__` STRING METADATA VIRTUAL, `__source__` STRING METADATA VIRTUAL, `__timestamp__` STRING METADATA VIRTUAL, __tag__ MAP<VARCHAR, VARCHAR> METADATA VIRTUAL, proctime as PROCTIME() ) WITH ( 'connector' = 'sls', 'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com', 'accessId' = '${secret_values.ak_id}', 'accessKey' = '${secret_values.ak_secret}', 'starttime' = '2023-08-30 00:00:00', 'project' ='sls-test', 'logstore' ='sls-input' ); CREATE TEMPORARY TABLE sls_sink( `time` BIGINT, url STRING, dt STRING, float_field FLOAT, double_field DOUBLE, boolean_field BOOLEAN, `__topic__` STRING, `__source__` STRING, `__timestamp__` BIGINT , receive_time BIGINT ) WITH ( 'connector' = 'sls', 'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com', 'accessId' = '${ak_id}', 'accessKey' = '${ak_secret}', 'project' ='sls-test', 'logstore' ='sls-output' ); INSERT INTO sls_sink SELECT `time`, url, dt, float_field, double_field, boolean_field, `__topic__` , `__source__` , `__timestamp__` , cast(__tag__['__receive_time__'] as bigint) as receive_time FROM sls_input; -
Ingestão de dados com uma source SLS
Use o SLS como fonte de dados para ingerir informações em tempo real em sistemas downstream suportados. Por exemplo, a configuração abaixo define um job de ingestão que grava dados de um logstore em um data lake formatado em Paimon no Data Lake Formation (DLF). O job infere automaticamente o schema da tabela de destino e suporta evolução de schema em tempo de execução.
source:
type: sls
name: SLS Source
endpoint: ${endpoint}
project: test_project
logstore: test_log
accessId: ${accessId}
accessKey: ${accessKey}
# Add a primary key to the table.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Route all data from test_project.test_log to the test_database.inventory table.
route:
- source-table: test_project.test_log
sink-table: test_database.inventory
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
DataStream API
Para ler ou gravar dados com a DataStream API, use um conector DataStream. Para mais informações, consulte Usage of DataStream connectors.
Se você usar uma versão do VVR anterior à 8.0.10, seu job pode falhar ao iniciar devido a dependências ausentes. Para resolver esse problema, adicione o uber-JAR correspondente como uma dependência adicional.
Leitura do SLS
O Realtime Compute for Apache Flink fornece a classe SlsSourceFunction, uma implementação de SourceFunction, para leitura de dados do Simple Log Service (SLS). O exemplo a seguir lê dados do SLS.
public class SlsDataStreamSource {
public static void main(String[] args) throws Exception {
// Sets up the streaming execution environment
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Creates an SLS source and prints the data to the console.
env.addSource(createSlsSource())
.map(SlsDataStreamSource::convertMessages)
.print();
env.execute("SLS Stream Source");
}
private static SlsSourceFunction createSlsSource() {
SLSAccessInfo accessInfo = new SLSAccessInfo();
accessInfo.setEndpoint("yourEndpoint");
accessInfo.setProjectName("yourProject");
accessInfo.setLogstore("yourLogStore");
accessInfo.setAccessId("yourAccessId");
accessInfo.setAccessKey("yourAccessKey");
// The batch get size is required.
accessInfo.setBatchGetSize(10);
// Optional parameters
accessInfo.setConsumerGroup("yourConsumerGroup");
accessInfo.setMaxRetries(3);
// Start time for consumption, set to the current time.
int startInSec = (int) (new Date().getTime() / 1000);
// Stop time for consumption, where -1 means never stop.
int stopInSec = -1;
return new SlsSourceFunction(accessInfo, startInSec, stopInSec);
}
private static List<String> convertMessages(SourceRecord input) {
List<String> res = new ArrayList<>();
for (FastLogGroup logGroup : input.getLogGroups()) {
int logsCount = logGroup.getLogsCount();
for (int i = 0; i < logsCount; i++) {
FastLog log = logGroup.getLogs(i);
int fieldCount = log.getContentsCount();
for (int idx = 0; idx < fieldCount; idx++) {
FastLogContent f = log.getContents(idx);
res.add(String.format("key: %s, value: %s", f.getKey(), f.getValue()));
}
}
}
return res;
}
}
Gravação no SLS
O Realtime Compute for Apache Flink fornece a classe SLSOutputFormat, uma implementação de OutputFormat, para gravação de dados no SLS. O exemplo a seguir grava dados no SLS.
public class SlsDataStreamSink {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.fromSequence(0, 100)
.map((MapFunction<Long, SinkRecord>) aLong -> getSinkRecord(aLong))
.addSink(createSlsSink())
.name(SlsDataStreamSink.class.getSimpleName());
env.execute("SLS Stream Sink");
}
private static OutputFormatSinkFunction createSlsSink() {
Configuration conf = new Configuration();
conf.setString(SLSOptions.ENDPOINT, "yourEndpoint");
conf.setString(SLSOptions.PROJECT, "yourProject");
conf.setString(SLSOptions.LOGSTORE, "yourLogStore");
conf.setString(SLSOptions.ACCESS_ID, "yourAccessId");
conf.setString(SLSOptions.ACCESS_KEY, "yourAccessKey");
SLSOutputFormat outputFormat = new SLSOutputFormat(conf);
return new OutputFormatSinkFunction<>(outputFormat);
}
private static SinkRecord getSinkRecord(Long seed) {
SinkRecord record = new SinkRecord();
LogItem logItem = new LogItem((int) (System.currentTimeMillis() / 1000));
logItem.PushBack("level", "info");
logItem.PushBack("name", String.valueOf(seed));
logItem.PushBack("message", "it's a test message for " + seed.toString());
record.setContent(logItem);
return record;
}
}
XML
O conector DataStream SLS está disponível no repositório central Maven.
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-sls</artifactId>
<version>${vvr-version}</version>
<exclusions>
<exclusion>
<groupId>org.apache.flink</groupId>
<artifactId>flink-format-common</artifactId>
</exclusion>
</exclusions>
</dependency>