O conector do DataHub permite ler dados em streaming do Alibaba Cloud DataHub em jobs do Flink e gravar os resultados processados de volta nos tópicos do DataHub. Ele oferece suporte tanto ao Flink SQL quanto à DataStream API.
O DataHub é compatível com o protocolo Kafka. Para conectar o Flink ao DataHub usando o protocolo Kafka, utilize o conector padrão do Kafka — não o conector Upsert Kafka. Para mais detalhes, consulte
.
Capacidades
|
Item |
Descrição |
|
Tipo suportado |
source e sink |
|
Modo de execução |
Streaming e batch |
|
Formato de dados |
N/A |
|
Métricas |
N/A |
|
Tipo de API |
DataStream e SQL |
|
Suporte a atualização/exclusão de dados no sink |
Não suportado. O sink grava apenas linhas de inserção no tópico de destino. |
Pré-requisitos
Antes de começar, verifique se você possui:
Um projeto e um tópico no DataHub. Consulte Introdução ao DataHub.
Uma assinatura do DataHub (necessária para o source). Consulte Criar uma assinatura.
Um AccessKey ID e um AccessKey secret da sua conta Alibaba Cloud. Consulte Operações no console.
Sintaxe
CREATE TEMPORARY TABLE datahub_input (
`time` BIGINT,
`sequence` STRING METADATA VIRTUAL,
`shard-id` BIGINT METADATA VIRTUAL,
`system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
'connector' = 'datahub',
'subId' = '<yourSubId>',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
Opções do conector
Geral
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
connector |
String |
Sim |
(nenhum) |
Tipo do conector. Defina como |
|
endPoint |
String |
Sim |
(nenhum) |
Endpoint do projeto DataHub. O valor varia conforme a região. Consulte Endpoints. |
|
project |
String |
Sim |
(nenhum) |
Nome do projeto DataHub. |
|
topic |
String |
Sim |
(nenhum) |
Nome do tópico DataHub. Para tópicos BLOB (dados não tipados e não estruturados), a tabela Flink deve conter exatamente uma coluna VARBINARY. |
|
accessId |
String |
Sim |
(nenhum) |
AccessKey ID da sua conta Alibaba Cloud. Armazene-o como variável em vez de codificá-lo diretamente. Consulte Gerenciar variáveis. |
|
accessKey |
String |
Sim |
(nenhum) |
AccessKey secret da sua conta Alibaba Cloud. |
|
retryTimeout |
Integer |
Não |
1800000 |
Tempo limite máximo para nova tentativa, em milissegundos. |
|
retryInterval |
Integer |
Não |
1000 |
Intervalo entre novas tentativas, em milissegundos. |
|
CompressType |
String |
Não |
lz4 |
Algoritmo de compactação para leituras e gravações. Valores válidos: |
Específico do source
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
subId |
String |
Sim |
(nenhum) |
ID da assinatura do DataHub. |
|
maxFetchSize |
Integer |
Não |
50 |
Quantidade de registros buscados por solicitação. Aumente este valor para melhorar o throughput de leitura. |
|
maxBufferSize |
Integer |
Não |
50 |
Número máximo de registros em cache de leituras assíncronas. Aumente este valor para melhorar o throughput de leitura. |
|
fetchLatestDelay |
Integer |
Não |
500 |
Tempo de espera em milissegundos quando não há dados disponíveis. Diminua este valor para reduzir a latência de leitura em tópicos com pouco tráfego. |
|
lengthCheck |
String |
Não |
NONE |
Regra para lidar com linhas em que a contagem de campos analisados não corresponde à contagem de colunas definidas. Valores válidos: |
|
columnErrorDebug |
Boolean |
Não |
false |
Define se os logs de depuração para erros de análise de campos devem ser ativados. Defina como |
|
startTime |
String |
Não |
(nenhum) |
Timestamp a partir do qual iniciar o consumo. Formato: |
|
endTime |
String |
Não |
(nenhum) |
Timestamp no qual parar o consumo. Formato: |
|
startTimeMs |
Long |
Não |
-1 |
Timestamp a partir do qual iniciar o consumo, em milissegundos. Tem precedência sobre |
Posição inicial de consumo
A opção startTimeMs controla onde o source começa a leitura:
-1 (padrão): Inicia a partir do offset mais recente no tópico. Se nenhum offset existir, retorna ao offset mais antigo.
Um timestamp específico: Inicia a partir do primeiro registro no timestamp especificado ou posterior a ele.
O valor padrão de -1 pode causar perda de dados. Se o job falhar antes do primeiro checkpoint, o offset mais recente no tópico pode ter avançado e os registros gravados durante essa janela serão ignorados. Defina startTimeMs explicitamente com um timestamp específico para controlar a posição inicial.
Regras de validação de contagem de campos
A opção lengthCheck determina o comportamento quando o número de campos analisados em uma linha não corresponde ao número de colunas definidas:
|
Valor |
Comportamento |
|
|
Se campos analisados > colunas definidas: lê da esquerda para a direita até a contagem definida. Se campos analisados < colunas definidas: ignora a linha. |
|
|
Ignora linhas em que a contagem de campos analisados difere da contagem de colunas definidas. |
|
|
Lança uma exceção quando a contagem de campos analisados difere da contagem de colunas definidas. |
|
|
Lê da esquerda para a direita. Se campos analisados > colunas definidas: lê até a contagem definida. Se campos analisados < colunas definidas: preenche os campos ausentes com |
Específico do sink
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
batchCount |
Integer |
Não |
500 |
Número máximo de linhas por lote de gravação. |
|
batchSize |
Integer |
Não |
512000 |
Tamanho máximo de um lote de gravação, em bytes. |
|
flushInterval |
Integer |
Não |
5000 |
Intervalo de liberação, em milissegundos. |
|
hashFields |
String |
Não |
null |
Lista separada por vírgulas de nomes de colunas usada para rotear linhas para shards. Linhas com os mesmos valores nessas colunas são gravadas no mesmo shard. O padrão ( |
|
timeZone |
String |
Não |
(nenhum) |
Fuso horário usado ao converter campos TIMESTAMP. |
|
schemaVersion |
Integer |
Não |
-1 |
Versão do esquema no registro de esquemas registrado. |
Comportamento de liberação do lote de gravação
O sistema libera um lote de gravação para o DataHub quando qualquer uma das seguintes condições é atendida primeiro:
O número de linhas no buffer atinge
batchCount.O tamanho total dos dados no buffer atinge
batchSize.O tempo desde a última liberação excede
flushInterval.
Aumentar batchCount, batchSize ou flushInterval melhora o throughput de gravação, mas aumenta a latência.
Mapeamentos de tipos de dados
|
Tipo Flink |
Tipo DataHub |
|
TINYINT |
TINYINT |
|
BOOLEAN |
BOOLEAN |
|
INTEGER |
INTEGER |
|
BIGINT |
BIGINT |
|
BIGINT |
TIMESTAMP |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
DECIMAL |
DECIMAL |
|
VARCHAR |
STRING |
|
SMALLINT |
SMALLINT |
|
VARBINARY |
BLOB |
Metadados
Os campos de metadados são somente leitura (R). Declare-os como METADATA VIRTUAL na definição da tabela source para incluí-los nas consultas sem gravá-los de volta no DataHub.
Os campos de metadados estão disponíveis apenas ao usar VVR 3.0.1 ou posterior.
|
Chave |
Tipo de dado |
Descrição |
R/W |
|
shard-id |
BIGINT METADATA VIRTUAL |
ID do shard do registro. |
R |
|
sequence |
STRING METADATA VIRTUAL |
Número de sequência do registro dentro do shard. |
R |
|
system-time |
TIMESTAMP METADATA VIRTUAL |
Momento em que o DataHub recebeu o registro. |
R |
Exemplos
Source
O exemplo a seguir lê dados de um tópico do DataHub e os imprime no console.
CREATE TEMPORARY TABLE datahub_input (
`time` BIGINT,
`sequence` STRING METADATA VIRTUAL,
`shard-id` BIGINT METADATA VIRTUAL,
`system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
'connector' = 'datahub',
'subId' = '<yourSubId>',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
CREATE TEMPORARY TABLE test_out (
`time` BIGINT,
`sequence` STRING,
`shard-id` BIGINT,
`system-time` TIMESTAMP
) WITH (
'connector' = 'print',
'logger' = 'true'
);
INSERT INTO test_out
SELECT
`time`,
`sequence`,
`shard-id`,
`system-time`
FROM datahub_input;
Sink
O exemplo abaixo lê de um tópico do DataHub, converte o campo name para minúsculas e grava os resultados em outro tópico do DataHub.
CREATE TEMPORARY TABLE datahub_source (
name VARCHAR
) WITH (
'connector' = 'datahub',
'endPoint' = '<endPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'subId' = '<yourSubId>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}',
'startTime' = '2018-06-01 00:00:00'
);
CREATE TEMPORARY TABLE datahub_sink (
name VARCHAR
) WITH (
'connector' = 'datahub',
'endPoint' = '<endPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}',
'batchSize' = '512000',
'batchCount' = '500'
);
INSERT INTO datahub_sink
SELECT
LOWER(name)
FROM datahub_source;
DataStream API
Para usar a DataStream API com o DataHub, configure um conector DataStream para o Realtime Compute for Apache Flink. Consulte Configurações de conectores DataStream.
Leitura do DataHub
O VVR fornece a classe DatahubSourceFunction, que implementa a interface SourceFunction do Flink.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// Configure the DataHub source
DatahubSourceFunction datahubSource =
new DatahubSourceFunction(
<yourEndPoint>,
<yourProjectName>,
<yourTopicName>,
<yourSubId>,
<yourAccessId>,
<yourAccessKey>,
"public",
<yourStartTime>,
<yourEndTime>
);
datahubSource.setRequestTimeout(30 * 1000);
datahubSource.enableExitAfterReadFinished();
env.addSource(datahubSource)
.map((MapFunction<RecordEntry, Tuple2<String, Long>>) this::getStringLongTuple2)
.print();
env.execute();
private Tuple2<String, Long> getStringLongTuple2(RecordEntry recordEntry) {
Tuple2<String, Long> tuple2 = new Tuple2<>();
TupleRecordData recordData = (TupleRecordData) (recordEntry.getRecordData());
tuple2.f0 = (String) recordData.getField(0);
tuple2.f1 = (Long) recordData.getField(1);
return tuple2;
}
Gravação no DataHub
O VVR fornece a classe OutputFormatSinkFunction, que implementa a interface DatahubSinkFunction.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Configure the DataHub sink
env.generateSequence(0, 100)
.map((MapFunction<Long, RecordEntry>) aLong -> getRecordEntry(aLong, "default:"))
.addSink(
new DatahubSinkFunction<>(
<yourEndPoint>,
<yourProjectName>,
<yourTopicName>,
<yourSubId>,
<yourAccessId>,
<yourAccessKey>,
"public",
<schemaVersion> // If schema registry is enabled, you must specify the schema version. Otherwise, set this to 0.
)
);
env.execute();
private RecordEntry getRecordEntry(Long message, String s) {
RecordSchema recordSchema = new RecordSchema();
recordSchema.addField(new Field("f1", FieldType.STRING));
recordSchema.addField(new Field("f2", FieldType.BIGINT));
recordSchema.addField(new Field("f3", FieldType.DOUBLE));
recordSchema.addField(new Field("f4", FieldType.BOOLEAN));
recordSchema.addField(new Field("f5", FieldType.TIMESTAMP));
recordSchema.addField(new Field("f6", FieldType.DECIMAL));
RecordEntry recordEntry = new RecordEntry();
TupleRecordData recordData = new TupleRecordData(recordSchema);
recordData.setField(0, s + message);
recordData.setField(1, message);
recordEntry.setRecordData(recordData);
return recordEntry;
}
Dependência Maven
Adicione o conector DataStream do DataHub ao seu projeto. Todas as versões disponíveis estão listadas no repositório central Maven.
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-datahub</artifactId>
<version>${vvr-version}</version>
</dependency>
Próximos passos
Para obter uma lista completa de conectores suportados pelo Realtime Compute for Apache Flink, consulte Conectores suportados.
Para se conectar ao DataHub usando o conector Kafka, consulte Message Queue for Apache Kafka.
Como retomar uma implantação que falha após a divisão ou redução de escala de um tópico do DataHub?
Posso excluir um tópico do DataHub que está sendo consumido?