Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector do DataHub

Última atualização: Jun 27, 2026

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.

Nota

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

Compatibilidade com Kafka

.

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:

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

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: lz4, deflate, "" (desativado). Requer VVR 6.0.5 ou posterior.

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: NONE, SKIP, EXCEPTION, PAD. Consulte Regras de validação de contagem de campos.

columnErrorDebug

Boolean

Não

false

Define se os logs de depuração para erros de análise de campos devem ser ativados. Defina como true para imprimir logs de exceção de análise.

startTime

String

Não

(nenhum)

Timestamp a partir do qual iniciar o consumo. Formato: yyyy-MM-dd hh:mm:ss.

endTime

String

Não

(nenhum)

Timestamp no qual parar o consumo. Formato: yyyy-MM-dd hh:mm:ss.

startTimeMs

Long

Não

-1

Timestamp a partir do qual iniciar o consumo, em milissegundos. Tem precedência sobre startTime. Consulte Posição inicial de consumo.

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.

Importante

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

NONE (padrão)

Se campos analisados > colunas definidas: lê da esquerda para a direita até a contagem definida. Se campos analisados < colunas definidas: ignora a linha.

SKIP

Ignora linhas em que a contagem de campos analisados difere da contagem de colunas definidas.

EXCEPTION

Lança uma exceção quando a contagem de campos analisados difere da contagem de colunas definidas.

PAD

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

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 (null) usa gravações aleatórias. Exemplo: hashFields=a,b.

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.

Nota

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

Importante

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