Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector do MaxCompute

Última atualização: Jun 27, 2026

O conector do MaxCompute permite ler e gravar dados no MaxCompute (anteriormente ODPS) — a plataforma de data warehouse totalmente gerenciada e em escala de exabytes da Alibaba Cloud — diretamente em jobs Flink SQL e DataStream.

Capacidades

Item

Descrição

Tipos de tabela

Tabela source, tabela de dimensão, tabela sink e sink de ingestão de dados

Modos de execução

Modo streaming e modo batch

Tipos de API

DataStream API, SQL API e jobs YAML de ingestão de dados

Semântica

At-least-once

Atualização ou exclusão de dados em uma tabela sink

Batch Tunnel ou Streaming Tunnel: apenas inserção. Upsert Tunnel: inserção, atualização e exclusão.

Métricas

Tipo de tabela

Métricas

Source

numRecordsIn, numRecordsInPerSecond, numBytesIn, numBytesInPerSecond

Sink

numRecordsOut, numRecordsOutPerSecond, numBytesOut, numBytesOutPerSecond

Tabela de dimensão

dim.odps.cacheSize

Para mais informações, consulte Monitoramento de métricas .

Pré-requisitos

Antes de começar, verifique se você tem:

Limitações

  • O conector suporta apenas a semântica at-least-once. Registros duplicados podem aparecer no MaxCompute dependendo do tunnel utilizado. Para obter orientações sobre a seleção de tunnels, consulte a seção "Como selecionar um tunnel de dados?" nas Perguntas frequentes sobre armazenamento upstream e downstream.

  • Por padrão, uma source opera em modo completo: ela lê apenas a partição especificada pela opção partition. Após a leitura de todos os dados, o job termina e não monitora novas partições. Para monitorar continuamente novas partições, configure uma source incremental usando startPartition.

  • Sempre que o cache de uma tabela de dimensão é atualizado, a tabela verifica a partição mais recente. Após o início da source, ela não lê dados recém-adicionados a uma partição que já está sendo lida. Execute um deployment somente depois que a partição contiver dados completos.

Escolha um tunnel

O MaxCompute oferece três tunnels para gravar dados a partir do Flink. Escolha com base no seu caso de uso:

Tunnel

Padrão

Quando usar

MaxCompute Batch Tunnel

Sim (useStreamTunnel=false, enableUpsert=false)

Cargas batch; os dados ficam disponíveis apenas após o checkpointing. Defina flushIntervalMs=0 para desativar o flush agendado.

MaxCompute Streaming Tunnel

Não (useStreamTunnel=true)

Ingestão near real-time; os dados submetidos a flush ficam imediatamente disponíveis no MaxCompute.

MaxCompute Upsert Tunnel

Não (enableUpsert=true)

Operações INSERT, UPDATE e DELETE em uma tabela Delta do MaxCompute. Requer VVR 8.0.6+.

Para uma comparação detalhada, consulte a seção "Como selecionar um tunnel de dados?" nas Perguntas frequentes sobre armazenamento upstream e downstream.

SQL

Use o conector do MaxCompute como tabela source, de dimensão ou sink em jobs baseados em SQL.

Sintaxe

CREATE TEMPORARY TABLE odps_source(
  id INT,
  user_name VARCHAR,
  content VARCHAR
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'schemaName' = '<yourSchemaName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=2018****'
);

Opções do conector

Opções gerais

Opção Obrigatória Padrão Tipo Descrição
connector Sim STRING Defina como odps.
endpoint Sim STRING Endpoint do MaxCompute. Consulte Endpoint.
tunnelEndpoint Não STRING Endpoint do MaxCompute Tunnel. Se não for especificado, o MaxCompute aloca conexões de tunnel por meio do Server Load Balancer (SLB).
project Sim STRING Nome do projeto do MaxCompute.
schemaName Não STRING Obrigatório apenas quando o recurso de schema do MaxCompute está ativado. Defina este valor como o nome do schema da tabela. Consulte Operações de schema. VVR 8.0.6+.
tableName Sim STRING Nome da tabela do MaxCompute.
accessId Sim STRING AccessKey ID usado para acessar o MaxCompute. Consulte Como visualizo meu AccessKey ID e AccessKey secret?
Importante

Armazene o AccessKey ID como uma variável. Consulte Gerenciar variáveis.

accessKey Sim STRING AccessKey secret usado para acessar o MaxCompute.
partition Não STRING Nome da partição na tabela do MaxCompute. Não é obrigatório para tabelas não particionadas ou sources incrementais. Consulte a seção "Como configuro a opção partition?" nas Perguntas frequentes sobre armazenamento upstream e downstream.
compressAlgorithm Não SNAPPY STRING Algoritmo de compressão para o MaxCompute Tunnel. Valores válidos: RAW (sem compressão), ZLIB, SNAPPY. Em cenários de teste, o SNAPPY melhora o throughput em aproximadamente 50% em comparação ao ZLIB.
quotaName Não STRING Nome da cota para grupos de recursos exclusivos do MaxCompute Tunnel. VVR 8.0.3+. Este parâmetro só tem efeito quando o endpoint está definido como um endereço VPC. Se um endpoint público for usado ou se o tunnelEndpoint for especificado, este parâmetro não terá efeito.

Opções de source

Opção

Obrigatória

Padrão

Tipo

Descrição

maxPartitionCount

Não

100

INTEGER

Número máximo de partições para leitura. Se excedido, o erro "The number of matched partitions exceeds the default limit" será exibido. Ler de muitas partições pode sobrecarregar o MaxCompute e tornar a inicialização do job mais lenta. Aumente este valor apenas quando sua carga de trabalho exigir.

useArrow

Não

false

BOOLEAN

Lê dados usando o formato Arrow, que chama a API de armazenamento do MaxCompute. Apenas deployments batch. VVR 8.0.8+.

splitSize

Não

256 MB

MEMORYSIZE

Quantidade de dados extraídos por split ao usar o formato Arrow. Apenas deployments batch. VVR 8.0.8+.

compressCodec

Não

"" (nenhum)

STRING

Algoritmo de compressão para leitura com o formato Arrow. Valores válidos: "" (nenhum), ZSTD, LZ4_FRAME. Especificar um codec melhora o throughput em comparação à ausência de compressão. Apenas deployments batch. VVR 8.0.8+.

dynamicLoadBalance

Não

false

BOOLEAN

Ativa a alocação dinâmica de shards para melhorar o desempenho de processamento e reduzir o tempo total de leitura. Observe que isso pode causar skew de dados, pois diferentes operadores leem quantidades inconsistentes de dados. Apenas deployments batch. VVR 8.0.8+.

Opções de source incremental

A source incremental consulta o MaxCompute intermitentemente para descobrir novas partições. Antes de ler uma nova partição, todas as gravações de dados nessa partição devem estar concluídas. Para detalhes, consulte a seção "O que faço se uma source incremental detectar uma nova partição enquanto os dados ainda estão sendo gravados?" nas Perguntas frequentes sobre armazenamento upstream e downstream.

Ordenação de partições: A source lê partições cuja ordem alfabética seja maior ou igual ao valor de startPartition. Por exemplo, year=2023,month=10 vem antes de year=2023,month=9 em ordem alfabética. Portanto, preencha os valores do mês com zero à esquerda (use year=2023,month=09 em vez de year=2023,month=9) para garantir a ordenação correta.

Opção

Obrigatória

Padrão

Tipo

Descrição

startPartition

Sim

STRING

Partição inicial para leituras incrementais. Quando especificada, partition é ignorada. Para tabelas com múltiplos níveis de particionamento, configure os valores das colunas de partição em ordem decrescente por nível. Consulte a seção "Como configuro startPartition?" nas Perguntas frequentes sobre armazenamento upstream e downstream.

subscribeIntervalInSec

Não

30

INTEGER

Intervalo de consulta em segundos.

modifiedTableOperation

Não

NONE

Enum

Ação quando uma partição é modificada durante a leitura. As sessões de download são salvas nos checkpoints; se os dados em uma partição forem alterados após o início de uma sessão, a retomada a partir do checkpoint falhará e o deployment reiniciará repetidamente. Valores válidos: NONE — atualiza startPartition para pular a partição indisponível e reinicia sem estado; SKIP — pula automaticamente a partição indisponível ao retomar. VVR 8.0.3+. Se definido como qualquer um dos valores, os dados já lidos da partição modificada são mantidos; os dados não lidos são descartados.

Opções de sink

Opção

Obrigatória

Padrão

Tipo

Descrição

useStreamTunnel

Não

false

BOOLEAN

Usa o MaxCompute Streaming Tunnel em vez do Batch Tunnel. true: Streaming Tunnel; false: Batch Tunnel. Consulte Escolha um tunnel.

flushIntervalMs

Não

30000 (30 s)

LONG

Intervalo de flush para o buffer do gravador de tunnel, em milissegundos. Para Streaming Tunnel: os dados submetidos a flush ficam imediatamente disponíveis. Para Batch Tunnel: os dados ficam disponíveis apenas após o checkpointing — defina como 0 para desativar o flush agendado. Acionado quando flushIntervalMs ou batchSize é atingido.

batchSize

Não

67108864 (64 MB)

LONG

Tamanho do buffer em bytes. Os dados são submetidos a flush quando o buffer atinge esse tamanho. Acionado quando batchSize ou flushIntervalMs é atingido.

numFlushThreads

Não

1

INTEGER

Número de threads usadas para submeter a flush o buffer do gravador de tunnel. Valores maiores que 1 permitem flush simultâneo entre partições.

slotNum

Não

0

INTEGER

Número de slots de Tunnel para receber dados do Flink. Consulte a Visão geral do serviço de transmissão de dados para limites de slots.

dynamicPartitionLimit

Não

100

INTEGER

Número máximo de partições dinâmicas gravadas entre dois checkpoints. Se excedido, o erro "Too many dynamic partitions" será exibido. Gravar em muitas partições aumenta a carga no MaxCompute e torna o checkpointing mais lento. Aumente este valor apenas quando sua carga de trabalho exigir.

retryTimes

Não

3

INTEGER

Número máximo de tentativas para solicitações ao servidor MaxCompute (criação de sessão, envio ou falhas de flush).

sleepMillis

Não

1000

INTEGER

Intervalo de nova tentativa em milissegundos.

enableUpsert

Não

false

BOOLEAN

Usa o MaxCompute Upsert Tunnel. true: processa registros INSERT, UPDATE_AFTER e DELETE; false: usa o tunnel especificado por useStreamTunnel. VVR 8.0.6+. Se o sink encontrar erros ou falhas de longa duração durante commits de sessão no modo upsert, defina o paralelismo do operador de sink para 10 ou menos.

upsertAsyncCommit

Não

false

BOOLEAN

Usa modo assíncrono ao fazer commit de sessões upsert. O modo assíncrono reduz o tempo de commit, mas os dados submetidos não ficam imediatamente consultáveis. VVR 8.0.6+.

upsertCommitTimeoutMs

Não

120000 (120 s)

INTEGER

Tempo limite para commits de sessão upsert, em milissegundos. VVR 8.0.6+.

sink.operation

Não

insert

STRING

Modo de gravação para uma tabela Delta. insert: modo de anexação; upsert: modo de atualização. VVR 8.0.10+.

sink.parallelism

Não

INTEGER

Paralelismo de gravação para uma tabela Delta. O padrão é o paralelismo upstream. O valor de write.bucket.num deve ser um múltiplo inteiro de sink.parallelism para obter desempenho ideal de gravação e eficiência de memória. VVR 8.0.10+.

sink.file-cached.enable

Não

false

BOOLEAN

Ativa o modo de cache de arquivo ao gravar em partições dinâmicas de uma tabela Delta. Reduz os arquivos pequenos gravados no servidor, mas aumenta a latência de gravação. Ative quando o sink tiver alto paralelismo. VVR 8.0.10+.

sink.file-cached.writer.num

Não

16

INTEGER

Threads de upload simultâneas por tarefa no modo de cache de arquivo. Evite definir um valor muito alto, pois gravar em muitas partições simultaneamente pode causar erros de falta de memória (OOM). Efetivo apenas quando sink.file-cached.enable=true. VVR 8.0.10+.

sink.bucket.check-interval

Não

60000

INTEGER

Intervalo de verificação de tamanho de arquivo no modo de cache de arquivo, em milissegundos. Efetivo apenas quando sink.file-cached.enable=true. VVR 8.0.10+.

sink.file-cached.rolling.max-size

Não

16 MB

MEMORYSIZE

Tamanho máximo de um único arquivo em cache. Quando excedido, os dados são carregados no servidor. Efetivo apenas quando sink.file-cached.enable=true. VVR 8.0.10+.

sink.file-cached.memory

Não

64 MB

MEMORYSIZE

Memória off-heap máxima para gravações de arquivo no modo de cache de arquivo. Efetivo apenas quando sink.file-cached.enable=true. VVR 8.0.10+.

sink.file-cached.memory.segment-size

Não

128 KB

MEMORYSIZE

Tamanho do segmento de buffer para gravações de arquivo no modo de cache de arquivo. Efetivo apenas quando sink.file-cached.enable=true. VVR 8.0.10+.

sink.file-cached.flush.always

Não

true

BOOLEAN

Se deve usar o cache ao gravar arquivos no modo de cache de arquivo. Efetivo apenas quando sink.file-cached.enable=true. VVR 8.0.10+.

sink.file-cached.write.max-retries

Não

3

INTEGER

Contagem de tentativas para uploads de dados no modo de cache de arquivo. Efetivo apenas quando sink.file-cached.enable=true. VVR 8.0.10+.

upsert.writer.max-retries

Não

3

INTEGER

Número máximo de tentativas para gravar em um bucket em uma sessão Upsert Writer. VVR 8.0.10+.

upsert.writer.buffer-size

Não

64 MB

MEMORYSIZE

Tamanho total do buffer em todos os buckets em uma sessão Upsert Writer. Os dados são submetidos a flush quando o total atinge esse limiar. Aumente para melhor eficiência de gravação; diminua se a gravação em muitas partições causar erros OOM. VVR 8.0.10+.

upsert.writer.bucket.buffer-size

Não

1 MB

MEMORYSIZE

Tamanho do buffer por bucket em uma sessão Upsert Writer. Diminua se a memória do servidor Flink for insuficiente. VVR 8.0.10+.

upsert.write.bucket.num

Sim

INTEGER

Número de buckets para a tabela Delta de destino. Deve corresponder ao write.bucket.num configurado na tabela Delta. VVR 8.0.10+.

upsert.write.slot-num

Não

1

INTEGER

Slots de Tunnel por sessão upsert. VVR 8.0.10+.

upsert.commit.max-retries

Não

3

INTEGER

Número máximo de tentativas para commits de sessão upsert. VVR 8.0.10+.

upsert.commit.thread-num

Não

16

INTEGER

Paralelismo para commits de sessão upsert. Evite valores grandes, pois commits simultâneos excessivos aumentam o consumo de recursos e podem causar problemas de desempenho. VVR 8.0.10+.

upsert.commit.timeout

Não

600

INTEGER

Tempo limite para commits de sessão upsert, em segundos. VVR 8.0.10+.

upsert.flush.concurrent

Não

2

INTEGER

Máximo de flushes simultâneos de bucket por partição. Cada flush de bucket ocupa um slot de Tunnel. VVR 8.0.10+.

insert.commit.thread-num

Não

16

INTEGER

Paralelismo para commits de sessão insert. VVR 8.0.10+.

insert.arrow-writer.enable

Não

false

BOOLEAN

Usa o formato Arrow para inserções. VVR 8.0.10+.

insert.arrow-writer.batch-size

Não

512

INTEGER

Máximo de linhas por lote no formato Arrow. VVR 8.0.10+.

insert.arrow-writer.flush-interval

Não

100000

INTEGER

Intervalo de flush do gravador em milissegundos. VVR 8.0.10+.

insert.writer.buffer-size

Não

64 MB

MEMORYSIZE

Tamanho do cache para o gravador com buffer. VVR 8.0.10+.

upsert.partial-column.enable

Não

false

BOOLEAN

Atualiza apenas colunas especificadas (atualização parcial de coluna). Aplica-se apenas a sinks de tabela Delta. Consulte Atualizar dados em colunas específicas. Quando true: se existir um registro com a mesma chave primária, os campos não nulos especificados serão sobrescritos; se não houver registro correspondente, um novo registro será inserido com novos valores para as colunas especificadas e null para todas as colunas não especificadas. VVR 8.0.11+.

Opções de tabela de dimensão

Quando um deployment é iniciado, a tabela de dimensão carrega todos os dados da partição especificada por partition. A opção partition suporta a função max_pt(). Ao recarregar o cache, a partição mais recente é relida. Defina partition como max_two_pt() para carregar dados de duas partições.

Tabelas de dimensão exigem cache=ALL . Aumente a memória do nó de join para pelo menos quatro vezes o tamanho dos dados da tabela remota. Para tabelas de dimensão grandes, use a dica SHUFFLE_HASH para distribuir os dados uniformemente. Para tabelas extremamente grandes que causam coletas de lixo (GCs) frequentes na Java Virtual Machine (JVM), converta para uma tabela de dimensão chave-valor com uma política de cache LRU (Least Recently Used) — por exemplo, uma tabela de dimensão do ApsaraDB for HBase.

Opção

Obrigatória

Padrão

Tipo

Descrição

cache

Sim

STRING

Política de cache. Deve ser definida como ALL e declarada explicitamente na instrução DDL. Todos os dados da tabela de dimensão são carregados no cache antes da execução do deployment. Consultas subsequentes pesquisam apenas no cache. O cache é recarregado após a expiração das entradas.

cacheSize

Não

100000

LONG

Número máximo de linhas a serem armazenadas em cache. Se excedido, o erro "Row count of table <table-name> partition <partition-name> exceeds maxRowCount limit" será exibido. Caches grandes consomem muita memória heap da JVM e tornam a inicialização e a atualização do cache mais lentas. Aumente este valor apenas quando sua carga de trabalho exigir.

cacheTTLMs

Não

Long.MAX_VALUE

LONG

Tempo limite do cache em milissegundos.

cacheReloadTimeBlackList

Não

STRING

Períodos durante os quais o cache não é atualizado. Use durante períodos de pico de tráfego (como eventos promocionais) para evitar instabilidade no deployment devido a atualizações de cache. Consulte a seção "Como configuro cacheReloadTimeBlackList?" nas Perguntas frequentes sobre armazenamento upstream e downstream.

maxLoadRetries

Não

10

INTEGER

Número máximo de tentativas para o carregamento inicial do cache na inicialização do deployment. Se as tentativas se esgotarem, o deployment falhará.

Mapeamentos de tipos de dados

Para a lista completa de tipos de dados do MaxCompute, consulte Sistema de tipos de dados do MaxCompute versão 2.0.

Tipo MaxCompute

Tipo Flink

BOOLEAN

BOOLEAN

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INTEGER

BIGINT

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

DECIMAL(precision, scale)

DECIMAL(precision, scale)

CHAR(n)

CHAR(n)

VARCHAR(n)

VARCHAR(n)

STRING

STRING

BINARY

BYTES

DATE

DATE

DATETIME

TIMESTAMP(3)

TIMESTAMP

TIMESTAMP(9)

TIMESTAMP_NTZ

TIMESTAMP(9)

ARRAY

ARRAY

MAP

MAP

STRUCT

ROW

JSON

STRING

Importante

Se uma tabela física do MaxCompute contiver campos de tipo composto aninhado (ARRAY, MAP ou STRUCT) ou um campo JSON, defina tblproperties('columnar.nested.type'='true') ao criar a tabela para permitir que o Realtime Compute for Apache Flink leia e grave dados corretamente.

Flink CDC (visualização pública)

O conector do MaxCompute pode ingerir dados de Change Data Capture (CDC) como sink em jobs baseados em YAML. Requer VVR 11.1+.

Sintaxe

source:
  type: xxx

sink:
  type: maxcompute
  name: MaxComputeSink
  access-id: ${your_accessId}
  access-key: ${your_accessKey}
  endpoint: ${your_maxcompute_endpoint}
  project: ${your_project}
  buckets-num: 8

Opções de configuração

Opção

Obrigatória

Padrão

Tipo

Descrição

type

Sim

String

Defina como maxcompute.

name

Não

String

Nome do sink.

access-id

Sim

String

AccessKey ID da sua conta Alibaba Cloud ou usuário RAM. Obtenha-o no console do Resource Access Management (RAM).

access-key

Sim

String

AccessKey secret.

endpoint

Sim

String

Endpoint do MaxCompute. Configure com base na região e no método de conexão de rede. Consulte Endpoint.

project

Sim

String

Nome do projeto do MaxCompute. Para encontrá-lo: faça login no console do MaxCompute, acesse Workspace > Projects e copie o nome do projeto.

tunnel.endpoint

Não

String

Endpoint do MaxCompute Tunnel. Geralmente inferido automaticamente. Obrigatório em ambientes de rede especiais, como com um servidor proxy.

quota.name

Não

String

Nome da cota para um grupo de recursos exclusivo. Se não for especificado, um grupo de recursos compartilhado será usado. Este parâmetro só tem efeito quando o endpoint está definido como um endereço VPC. Se um endpoint público for usado ou se o tunnel.endpoint for especificado, este parâmetro não terá efeito.

sts-token

Não

String

Token do Security Token Service (STS) para autenticação de função RAM. Obrigatório ao acessar o MaxCompute com uma função RAM.

buckets-num

Não

16

Integer

Número de buckets para uma tabela Delta do MaxCompute criada automaticamente. Consulte Data warehouse near real-time.

compress.algorithm

Não

zlib

String

Algoritmo de compressão de dados. Valores válidos: raw (sem compressão), zlib, snappy.

total.buffer-size

Não

64 MB

String

Tamanho do buffer na memória. Para tabelas particionadas: aplica-se por partição. Para tabelas não particionadas: aplica-se por tabela. Buffers para diferentes partições ou tabelas são independentes. Os dados são submetidos a flush quando o buffer está cheio.

bucket.buffer-size

Não

4 MB

String

Tamanho do buffer por bucket. Aplica-se apenas ao gravar em tabelas Delta do MaxCompute.

commit.thread-num

Não

16

Integer

Máximo de partições ou tabelas submetidas a commit simultaneamente durante o checkpointing.

flush.concurrent-num

Não

4

Integer

Máximo de buckets submetidos a flush simultaneamente. Aplica-se apenas ao gravar em tabelas Delta do MaxCompute.

Mapeamentos de localização de tabela

Quando o conector cria tabelas automaticamente no MaxCompute, as localizações são mapeadas da seguinte forma:

Importante

Se o recurso de schema estiver desativado para seu projeto do MaxCompute, o conector ignora tableId.namespace. Nesse caso, apenas um único banco de dados (ou seu equivalente lógico) é ingerido no MaxCompute — por exemplo, apenas um banco de dados MySQL ao ingerir do MySQL.

Localização MySQL

Abstração Flink CDC

Localização MaxCompute

N/A

Projeto (da configuração)

Projeto

Banco de dados

TableId.namespace

Schema (ignorado se o schema estiver desativado)

Tabela

TableId.tableName

Tabela

Mapeamentos de tipos de dados

Tipo Flink CDC

Tipo MaxCompute

CHAR

STRING

VARCHAR

STRING

BOOLEAN

BOOLEAN

BINARY/VARBINARY

BINARY

DECIMAL

DECIMAL

TINYINT

TINYINT

SMALLINT

SMALLINT

INTEGER

INTEGER

BIGINT

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

TIME_WITHOUT_TIME_ZONE

STRING

DATE

DATE

TIMESTAMP_WITHOUT_TIME_ZONE

TIMESTAMP_NTZ

TIMESTAMP_WITH_LOCAL_TIME_ZONE (precision > 3)

TIMESTAMP

TIMESTAMP_WITH_LOCAL_TIME_ZONE (precision <= 3)

DATETIME

TIMESTAMP_WITH_TIME_ZONE (precision > 3)

TIMESTAMP

TIMESTAMP_WITH_TIME_ZONE (precision <= 3)

DATETIME

ARRAY

ARRAY

MAP

MAP

ROW

STRUCT

Exemplos

SQL API

Tabela source

Ler todos os dados de uma partição

Leia todos os dados da partição especificada por partition:

CREATE TEMPORARY TABLE odps_source (
  cid VARCHAR,
  rt DOUBLE
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpointName>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=201809*'
);

CREATE TEMPORARY TABLE blackhole_sink (
  cid VARCHAR,
  invoke_count BIGINT
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT
   cid,
   COUNT(*) AS invoke_count
FROM odps_source GROUP BY cid;

Ler dados incrementais

Leia dados a partir da partição especificada por startPartition e monitore continuamente novas partições:

CREATE TEMPORARY TABLE odps_source (
  cid VARCHAR,
  rt DOUBLE
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpointName>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'startPartition' = 'yyyy=2018,MM=09,dd=05' -- Start reading from the 20180905 partition.
);

CREATE TEMPORARY TABLE blackhole_sink (
  cid VARCHAR,
  invoke_count BIGINT
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT cid, COUNT(*) AS invoke_count
FROM odps_source GROUP BY cid;

Tabela sink

Gravar em uma partição estática

Grave na partição especificada por partition:

CREATE TEMPORARY TABLE datagen_source (
  id INT,
  len INT,
  content VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_sink (
  id INT,
  len INT,
  content VARCHAR
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905' -- Write to partition 20180905.
);

INSERT INTO odps_sink
SELECT
  id, len, content
FROM datagen_source;

Gravar em partições dinâmicas

Grave dados em partições determinadas em tempo de execução pelos valores na coluna ds:

CREATE TEMPORARY TABLE datagen_source (
  id INT,
  len INT,
  content VARCHAR,
  c TIMESTAMP
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_sink (
  id  INT,
  len INT,
  content VARCHAR,
  ds VARCHAR -- Dynamic partition column.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds' -- Omit the value; data is routed to partitions based on the ds field.
);

INSERT INTO odps_sink
SELECT
   id,
   len,
   content,
   DATE_FORMAT(c, 'yyMMdd') as ds
FROM datagen_source;

Tabela de dimensão

Chave de valor único

Especifique uma chave primária quando cada chave mapear exatamente uma linha:

CREATE TEMPORARY TABLE datagen_source (
  k INT,
  v VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_dim (
  k INT,
  v VARCHAR,
  PRIMARY KEY (k) NOT ENFORCED  -- Specify the primary key.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905',
  'cache' = 'ALL'
);

CREATE TEMPORARY TABLE blackhole_sink (
  k VARCHAR,
  v1 VARCHAR,
  v2 VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT k, s.v, d.v
FROM datagen_source AS s
INNER JOIN odps_dim FOR SYSTEM_TIME AS OF PROCTIME() AS d ON s.k = d.k;

Chave de múltiplos valores

Omita a chave primária quando uma chave puder mapear várias linhas:

CREATE TEMPORARY TABLE datagen_source (
  k INT,
  v VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_dim (
  k INT,
  v VARCHAR
  -- No primary key needed for multi-value lookups.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905',
  'cache' = 'ALL'
);

CREATE TEMPORARY TABLE blackhole_sink (
  k VARCHAR,
  v1 VARCHAR,
  v2 VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT k, s.v, d.v
FROM datagen_source AS s
INNER JOIN odps_dim FOR SYSTEM_TIME AS OF PROCTIME() AS d ON s.k = d.k;

DataStream API

Importante
  • Para usar a DataStream API com o MaxCompute, configure um conector DataStream. Consulte Integrar conectores DataStream.

  • O VVR 6.0.6+ suporta depuração local de programas DataStream com o conector do MaxCompute por até 30 minutos. Sessões que excederem 30 minutos são encerradas com um erro. Consulte Depurar conectores localmente.

  • A leitura de uma tabela Delta do MaxCompute (uma tabela criada com uma chave primária e transactional=true) não é suportada.

Declare a tabela do MaxCompute usando SQL e, em seguida, acesse-a por meio da Table API ou DataStream API.

Conectar à tabela source

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
tEnv.executeSql(String.join(
    "\n",
    "CREATE TEMPORARY TABLE IF NOT EXISTS odps_source (",
    "  cid VARCHAR,",
    "  rt DOUBLE",
    ") WITH (",
    "  'connector' = 'odps',",
    "  'endpoint' = '<yourEndpointName>',",
    "  'project' = '<yourProjectName>',",
    "  'tableName' = '<yourTableName>',",
    "  'accessId' = '<yourAccessId>',",
    "  'accessKey' = '<yourAccessPassword>',",
    "  'partition' = 'ds=201809*'",
    ")");
DataStream<Row> source = tEnv.toDataStream(tEnv.from("odps_source"));
source.print();
env.execute("odps source");

Conectar ao sink

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
tEnv.executeSql(String.join(
    "\n",
    "CREATE TEMPORARY TABLE IF NOT EXISTS odps_sink (",
    "  cid VARCHAR,",
    "  rt DOUBLE",
    ") WITH (",
    "  'connector' = 'odps',",
    "  'endpoint' = '<yourEndpointName>',",
    "  'project' = '<yourProjectName>',",
    "  'tableName' = '<yourTableName>',",
    "  'accessId' = '<yourAccessId>',",
    "  'accessKey' = '<yourAccessPassword>',",
    "  'partition' = 'ds=20180905'",
    ")");
DataStream<Row> data = env.fromElements(
    Row.of("id0", 3.),
    Row.of("id1", 4.));
tEnv.fromDataStream(data).insertInto("odps_sink").execute();

Dependência Maven

Adicione o conector DataStream do MaxCompute ao seu projeto. Todas as versões estão disponíveis no repositório central Maven.

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-odps</artifactId>
    <version>${vvr-version}</version>
</dependency>

Próximos passos