O VVR 11.6.0 adiciona funções de IA multimodal (PDFs e imagens), o tipo Variant para dados semiestruturados e promove a Ingestão de Dados (Flink CDC) para disponibilidade geral (GA). Esta versão também aprimora os conectores Kafka, MySQL CDC, OceanBase, Elasticsearch e Hologres, além de incorporar correções das comunidades Apache Flink 1.20.2 e 1.20.3.
Esta atualização está sendo implementada em fases em todas as regiões. Se sua conta ainda não foi atualizada, esses recursos estarão indisponíveis. Para solicitar uma atualização antecipada, envie um ticket.
Aprimoramentos do mecanismo
IA
Novas funções integradas para processamento de imagens e documentos em tempo real:
Conversão de PDF para imagem: converta páginas de PDF em imagens dentro de um job do Flink.
Recuperação de arquivos: leia o conteúdo de arquivos diretamente do OSS e do Message Notification Service (MNS).
Detecção de nitidez de imagem: avalie a nitidez da imagem com OpenCV.
Compressão de imagem: reduza o tamanho da imagem inline.
Repasse de imagem Base64: encaminhe imagens codificadas em Base64 para chamadas de modelos downstream.
Essas funções oferecem suporte a pipelines de inferência ponta a ponta com Vision Language Models (VLMs), como o Qwen-VL.
Conector Simple Message Queue (MNS)
Um novo conector MNS permite assinar eventos de alteração do OSS e fechar o ciclo em pipelines multimodais em tempo real.
Flink SQL
Tipo Variant
O novo tipo VARIANT lida com dados semiestruturados sem esquema fixo. Acesse campos com notação de ponto (variant.field) ou colchetes (variant['key']), converta tipos básicos e grave colunas Variant em sinks Paimon.
Novas funções
parse_json: converte uma string JSON em um valor Variant (disponível noTransformdo Flink CDC).MD5 e outras funções de hash: disponíveis no
Transformdo Flink CDC.
Ingestão de Dados (Flink CDC) — agora GA
A Ingestão de Dados (Flink CDC) saiu da prévia pública e agora tem disponibilidade geral (GA).
Status de GA dos conectores
|
Status |
Conectores |
|
Disponibilidade geral |
Paimon, StarRocks, Hologres, MySQL, Kafka |
|
Prévia pública |
Doris, OceanBase, MaxCompute, SLS, MongoDB, Postgres |
Novidades nesta versão
Mesclagem de colunas
Combine vários campos upstream com nomes ou capitalizações diferentes em uma única coluna de destino. Esse recurso oferece suporte a correspondência por regex, normalização de maiúsculas/minúsculas e mapeamentos personalizados, sendo útil para consolidar fontes JSON com nomenclatura de campos inconsistente.
Gravações apenas de acréscimo em tabelas particionadas do Paimon
O sink Paimon agora grava em tabelas particionadas sem chave primária. Anteriormente, a chave de partição precisava fazer parte da chave primária; esse requisito foi removido para casos de uso apenas de acréscimo.
Se você depende atualmente do comportamento antigo, verifique a configuração do seu sink Paimon após a atualização.
Aprimoramentos do Transform
Limpe uma chave primária ou de partição passando um valor
null.Defina roteamentos complexos de nomes de tabela usando expressões regulares.
Suporte a Variant
Acesse e converta campos do tipo
VARIANT; grave dados Variant no Paimon.
Aprimoramentos de source
Source PolarDB-X CDC (prévia pública): assinatura de binlog de alto paralelismo no nível da tabela e assinatura por dimensão de tabela.
Source SLS: imponha tipos específicos de análise de campos.
Source Kafka: divida uma única mensagem do Kafka em vários registros com base nos valores dos campos (roteamento de campo) e grave-os em diferentes tabelas de destino. Particionadores personalizados também são suportados.
Aprimoramentos de sink
|
Sink |
Aprimoramento |
|
Paimon |
Configuração separada de paralelismo para nós de commit |
|
MaxCompute |
Mapeamento de tipo DATETIME; lógica de commit otimizada para reduzir o consumo de consultas por segundo (QPS) |
|
Iceberg |
Referências de catálogo integradas; recuperação automática de URLs de conexão e credenciais |
Conectores
Kafka
O sink grava IDs de tabela de três partes (
Database.Schema.Table) no formato Debezium JSON.Após uma alteração de tópico, uma reinicialização com estado agora gera uma exceção de estado incompatível em vez de consumir silenciosamente dos tópicos antigo e novo simultaneamente. Isso evita consumo duplo e inconsistências de dados.
Se seus jobs dependem de reinicialização após uma mudança de tópico, revise esse comportamento antes de atualizar.
MySQL CDC
As mensagens de erro para identificadores globais de transação (GTIDs) expirados agora indicam claramente a causa raiz.
O ID do servidor consumidor é incluído nos logs para simplificar a solução de problemas.
PolarDB-X
Agora há suporte como source YAML CDC (prévia pública).
OceanBase
O sink JDBC agora oferece suporte a rollbacks manuais de transação e reutilização de pool de conexões, resolvendo desconexões frequentes causadas por wait_timeout.
Elasticsearch
A tabela source e a tabela de dimensão agora oferecem suporte ao Elasticsearch 8.x (com o cliente compatível com ES7).
Doris
As mensagens de erro para configurações incorretas de porta estão mais descritivas.
Integração com data lakehouse
Iceberg
O sink agora oferece suporte à métrica
numRecordsOutOfSinkPerSecond(OUT RPS).Configure parâmetros relacionados ao Hadoop para maior flexibilidade de conexão.
Jobs do Flink CDC podem gravar no Iceberg do Data Lake Formation (DLF).
Hologres
A tabela source Binlog permite consumo a partir do offset
LATEST.O catálogo Hologres expõe índices secundários e chaves de varredura por prefixo.
A leitura do tipo de array
varchar[]agora é suportada.O cache de detecção de parâmetros foi otimizado para evitar tempos limite de inicialização quando há muitas tabelas.
O sink aceita paralelismo superior ao número de shards quando
sink.reshuffle-by-holo-distribution-key.enabledestá configurado.
MaxCompute
O catálogo usa consultas paginadas para evitar o congelamento do centro de metadados.
A lógica de commit do sink YAML otimizada reduz erros OOM causados por limites de QPS.
Hive
O catálogo permite especificar um formato de armazenamento (como Parquet) ao criar uma tabela.
Paimon
Adicionado suporte ao formato de arquivo Lance.
Observabilidade
Novas métricas:
|
Métrica |
Descrição |
|
|
Uso de disco local |
|
|
Memória nativa do GeminiDB em uso |
|
|
Limite de memória nativa do GeminiDB |
|
|
Limite superior de paralelismo do operador auto-pilot |
Logs WARN não essenciais — como alertas quando um formato não oferece suporte a snapshots — agora são suprimidos.
Correções de bugs
Estabilidade
Correções da comunidade Apache Flink: incorpora correções das versões 1.20.2 e 1.20.3 do Apache Flink.
Perda de dados no Kafka: corrigida a perda de dados ao ler do Kafka e gravar no OSS com transações ativadas.
Desconexão do PolarDB-X: corrigido um pico abrupto de latência e a exceção
EOFExceptionacionados por uma conexão interrompida do PolarDB-X.wait_timeoutdo OceanBase: corrigidas desconexões frequentes no sink JDBC do OceanBase causadas porwait_timeout.
Corretude
Canal Protobuf do Flink CDC: corrigido o formato de timestamp inconsistente e o tratamento do tipo
tinyint.Depuração do source MySQL CDC: corrigida uma contagem com erro de um dígito quando a reutilização estava ativada (a contagem de tabelas agora exibe o número correto).
Sink MaxCompute (ODPS) do Flink CDC: corrigidos erros OOM de Metaspace causados por commits frequentes.
Experiência
Mensagens de erro de Temporal Join: as mensagens de erro para sintaxe de Temporal Join estão mais claras.
Logs WARN internos:
Cannot snapshot the tablee alertas internos semelhantes agora são registrados no nível DEBUG.Campos nulos no binlog do Hologres: corrigida uma exceção que ocorria quando alguns campos eram nulos durante o consumo de binlog do Hologres.