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. |
Pré-requisitos
Antes de começar, verifique se você tem:
Uma tabela do MaxCompute. Consulte Criar uma tabela.
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 usandostartPartition.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 ( |
Cargas batch; os dados ficam disponíveis apenas após o checkpointing. Defina |
|
MaxCompute Streaming Tunnel |
Não ( |
Ingestão near real-time; os dados submetidos a flush ficam imediatamente disponíveis no MaxCompute. |
|
MaxCompute Upsert Tunnel |
Não ( |
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 |
|
|
Não |
|
INTEGER |
Número máximo de partições para leitura. Se excedido, o erro |
|
|
Não |
|
BOOLEAN |
Lê dados usando o formato Arrow, que chama a API de armazenamento do MaxCompute. Apenas deployments batch. VVR 8.0.8+. |
|
|
Não |
|
MEMORYSIZE |
Quantidade de dados extraídos por split ao usar o formato Arrow. Apenas deployments batch. VVR 8.0.8+. |
|
|
Não |
|
STRING |
Algoritmo de compressão para leitura com o formato Arrow. Valores válidos: |
|
|
Não |
|
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 |
|
|
Sim |
— |
STRING |
Partição inicial para leituras incrementais. Quando especificada, |
|
|
Não |
|
INTEGER |
Intervalo de consulta em segundos. |
|
|
Não |
|
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: |
Opções de sink
|
Opção |
Obrigatória |
Padrão |
Tipo |
Descrição |
|
|
Não |
|
BOOLEAN |
Usa o MaxCompute Streaming Tunnel em vez do Batch Tunnel. |
|
|
Não |
|
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 |
|
|
Não |
|
LONG |
Tamanho do buffer em bytes. Os dados são submetidos a flush quando o buffer atinge esse tamanho. Acionado quando |
|
|
Não |
|
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. |
|
|
Não |
|
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. |
|
|
Não |
|
INTEGER |
Número máximo de partições dinâmicas gravadas entre dois checkpoints. Se excedido, o erro |
|
|
Não |
|
INTEGER |
Número máximo de tentativas para solicitações ao servidor MaxCompute (criação de sessão, envio ou falhas de flush). |
|
|
Não |
|
INTEGER |
Intervalo de nova tentativa em milissegundos. |
|
|
Não |
|
BOOLEAN |
Usa o MaxCompute Upsert Tunnel. |
|
|
Não |
|
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+. |
|
|
Não |
|
INTEGER |
Tempo limite para commits de sessão upsert, em milissegundos. VVR 8.0.6+. |
|
|
Não |
|
STRING |
Modo de gravação para uma tabela Delta. |
|
|
Não |
— |
INTEGER |
Paralelismo de gravação para uma tabela Delta. O padrão é o paralelismo upstream. O valor de |
|
|
Não |
|
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+. |
|
|
Não |
|
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 |
|
|
Não |
|
INTEGER |
Intervalo de verificação de tamanho de arquivo no modo de cache de arquivo, em milissegundos. Efetivo apenas quando |
|
|
Não |
|
MEMORYSIZE |
Tamanho máximo de um único arquivo em cache. Quando excedido, os dados são carregados no servidor. Efetivo apenas quando |
|
|
Não |
|
MEMORYSIZE |
Memória off-heap máxima para gravações de arquivo no modo de cache de arquivo. Efetivo apenas quando |
|
|
Não |
|
MEMORYSIZE |
Tamanho do segmento de buffer para gravações de arquivo no modo de cache de arquivo. Efetivo apenas quando |
|
|
Não |
|
BOOLEAN |
Se deve usar o cache ao gravar arquivos no modo de cache de arquivo. Efetivo apenas quando |
|
|
Não |
|
INTEGER |
Contagem de tentativas para uploads de dados no modo de cache de arquivo. Efetivo apenas quando |
|
|
Não |
|
INTEGER |
Número máximo de tentativas para gravar em um bucket em uma sessão Upsert Writer. VVR 8.0.10+. |
|
|
Não |
|
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+. |
|
|
Não |
|
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+. |
|
|
Sim |
— |
INTEGER |
Número de buckets para a tabela Delta de destino. Deve corresponder ao |
|
|
Não |
|
INTEGER |
Slots de Tunnel por sessão upsert. VVR 8.0.10+. |
|
|
Não |
|
INTEGER |
Número máximo de tentativas para commits de sessão upsert. VVR 8.0.10+. |
|
|
Não |
|
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+. |
|
|
Não |
|
INTEGER |
Tempo limite para commits de sessão upsert, em segundos. VVR 8.0.10+. |
|
|
Não |
|
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+. |
|
|
Não |
|
INTEGER |
Paralelismo para commits de sessão insert. VVR 8.0.10+. |
|
|
Não |
|
BOOLEAN |
Usa o formato Arrow para inserções. VVR 8.0.10+. |
|
|
Não |
|
INTEGER |
Máximo de linhas por lote no formato Arrow. VVR 8.0.10+. |
|
|
Não |
|
INTEGER |
Intervalo de flush do gravador em milissegundos. VVR 8.0.10+. |
|
|
Não |
|
MEMORYSIZE |
Tamanho do cache para o gravador com buffer. VVR 8.0.10+. |
|
|
Não |
|
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 |
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 exigemcache=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 dicaSHUFFLE_HASHpara 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 |
|
|
Sim |
— |
STRING |
Política de cache. Deve ser definida como |
|
|
Não |
|
LONG |
Número máximo de linhas a serem armazenadas em cache. Se excedido, o erro |
|
|
Não |
|
LONG |
Tempo limite do cache em milissegundos. |
|
|
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. |
|
|
Não |
|
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 |
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 |
|
|
Sim |
— |
String |
Defina como |
|
|
Não |
— |
String |
Nome do sink. |
|
|
Sim |
— |
String |
AccessKey ID da sua conta Alibaba Cloud ou usuário RAM. Obtenha-o no console do Resource Access Management (RAM). |
|
|
Sim |
— |
String |
AccessKey secret. |
|
|
Sim |
— |
String |
Endpoint do MaxCompute. Configure com base na região e no método de conexão de rede. Consulte Endpoint. |
|
|
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. |
|
|
Não |
— |
String |
Endpoint do MaxCompute Tunnel. Geralmente inferido automaticamente. Obrigatório em ambientes de rede especiais, como com um servidor proxy. |
|
|
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 |
|
|
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. |
|
|
Não |
|
Integer |
Número de buckets para uma tabela Delta do MaxCompute criada automaticamente. Consulte Data warehouse near real-time. |
|
|
Não |
|
String |
Algoritmo de compressão de dados. Valores válidos: |
|
|
Não |
|
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. |
|
|
Não |
|
String |
Tamanho do buffer por bucket. Aplica-se apenas ao gravar em tabelas Delta do MaxCompute. |
|
|
Não |
|
Integer |
Máximo de partições ou tabelas submetidas a commit simultaneamente durante o checkpointing. |
|
|
Não |
|
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:
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 |
|
Schema (ignorado se o schema estiver desativado) |
|
Tabela |
|
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
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>