Os ClickPipes do Kinesis podem ser implantados e gerenciados manualmente por meio da UI do ClickPipes, bem como programaticamente usando OpenAPI e Terraform.
Você já se familiarizou com a introdução ao ClickPipes e configurou credenciais do IAM ou uma função do IAM. Siga o guia de acesso baseado em função do Kinesis para saber como configurar uma função que funcione com o ClickHouse Cloud.
Criando seu primeiro ClickPipe
- Acesse o SQL Console do seu serviço ClickHouse Cloud.
- Selecione o botão
Data Sources no menu à esquerda e clique em “Set up a ClickPipe”
- Selecione sua fonte de dados.
- Preencha o formulário informando um nome para o ClickPipe, uma descrição (opcional), sua função do IAM ou credenciais e outros detalhes da conexão.
- Selecione o Kinesis Stream e o offset inicial. A UI exibirá um documento de exemplo da fonte selecionada (Kafka topic etc.). Você também pode ativar o Enhanced Fan-out para streams do Kinesis para melhorar o desempenho e a estabilidade do seu ClickPipe (mais informações sobre o Enhanced Fan-out podem ser encontradas aqui)
- Na próxima etapa, você pode escolher se deseja fazer a ingestão de dados em uma nova tabela do ClickHouse ou reutilizar uma existente. Siga as instruções na tela para modificar o nome da tabela, o esquema e as configurações. Você pode ver uma prévia em tempo real das suas alterações na tabela de exemplo na parte superior.
Você também pode personalizar as configurações avançadas usando os controles fornecidos
- Como alternativa, você pode optar por fazer a ingestão dos seus dados em uma tabela existente do ClickHouse. Nesse caso, a UI permitirá mapear campos da fonte para os campos do ClickHouse na tabela de destino selecionada.
- Por fim, você pode configurar as permissões para o usuário interno do ClickPipes.
Permissões: o ClickPipes criará um usuário dedicado para gravar dados em uma tabela de destino. Você pode selecionar uma função para esse usuário interno usando uma função personalizada ou uma das funções predefinidas:
Full access: com acesso total ao cluster. Isso pode ser útil se você usar visão materializada ou Dicionário com a tabela de destino.
Only destination table: apenas com permissões INSERT na tabela de destino.
- Ao clicar em “Complete Setup”, o sistema registrará seu ClickPipe, e você poderá vê-lo listado na tabela de resumo.
A tabela de resumo fornece controles para exibir dados de exemplo da fonte ou da tabela de destino no ClickHouse
Bem como controles para remover o ClickPipe e exibir um resumo do job de ingestão.
- Parabéns! você configurou com sucesso seu primeiro ClickPipe. Se este for um ClickPipe de streaming, ele será executado continuamente, fazendo ingestão de dados em tempo real da sua fonte de dados remota. Caso contrário, ele fará a ingestão do lote e será concluído.
Os formatos suportados são:
O ClickPipes para Kinesis detecta e descompacta automaticamente registros compactados. Diferentemente do Kafka, em que a biblioteca cliente cuida da descompactação de forma transparente, o Kinesis entrega bytes brutos — o ClickPipes faz isso para você, sem exigir nenhuma configuração.
Os seguintes codecs de compressão são suportados:
- gzip
- zstd
- lz4
- snappy (formato com frames)
A compressão é detectada automaticamente por meio dos bytes mágicos de cada registro. Se nenhuma assinatura de compressão conhecida for encontrada, o registro será tratado como sem compactação. O tipo de compressão detectado também é exibido durante a inferência de esquema, para que a prévia dos dados de amostra na UI mostre corretamente os dados descompactados.
A detecção automática é segura para formatos baseados em texto, como JSON e CSV, pois caracteres ASCII imprimíveis nunca coincidem com bytes mágicos de compressão.
Tipos de dados suportados
Os seguintes tipos de dados do ClickHouse são compatíveis no momento com o ClickPipes:
- Tipos numéricos básicos - [U]Int8/16/32/64, Float32/64 e BFloat16
- Tipos inteiros grandes - [U]Int128/256
- Tipos Decimal
- Boolean
- String
- FixedString
- Date, Date32
- DateTime, DateTime64 (apenas timezones UTC)
- Enum8/Enum16
- UUID
- IPv4
- IPv6
- todos os tipos LowCardinality do ClickHouse
- map com chaves e valores usando qualquer um dos tipos acima (incluindo Nullable)
- Tuple e Array com elementos usando qualquer um dos tipos acima (incluindo Nullable, com apenas um nível de profundidade)
- Tipos SimpleAggregateFunction (para destinations AggregatingMergeTree ou SummingMergeTree)
Você pode especificar manualmente um tipo Variant (como Variant(String, Int64, DateTime)) para qualquer campo JSON
no stream de dados de origem. Devido à forma como o ClickPipes determina o subtipo correto de Variant a ser usado, apenas um tipo inteiro ou datetime
pode ser usado na definição de Variant — por exemplo, Variant(Int64, UInt32) não é compatível.
Campos JSON que são sempre um objeto JSON podem ser atribuídos a uma coluna de destino do tipo JSON. Você terá que alterar manualmente a
coluna de destino para o tipo JSON desejado, incluindo quaisquer caminhos fixos ou ignorados.
Colunas virtuais do Kinesis
As colunas virtuais a seguir são compatíveis com o stream do Kinesis. Ao criar uma nova tabela de destino, é possível adicionar colunas virtuais usando o botão Add Column.
O campo _raw_message pode ser usado nos casos em que apenas o registro JSON completo do Kinesis é necessário (como ao usar as funções JsonExtract* do ClickHouse para preencher uma
visão materializada downstream). Para esses pipes, excluir todas as colunas “não virtuais” pode melhorar o desempenho do ClickPipes.
- DEFAULT não é suportado.
- Mensagens individuais são limitadas por padrão a 16 MB (sem compactação) ao usar o menor tamanho de réplica (XS) e a 32 MB (sem compactação) com réplicas maiores. Mensagens que excederem esse limite serão rejeitadas com erro. Se precisar de mensagens maiores, entre em contato com o suporte.
O ClickPipes insere dados no ClickHouse em lotes. Isso evita a criação de um número excessivo de partes no banco de dados, o que pode levar a problemas de desempenho no cluster.
Os lotes são inseridos quando um dos seguintes critérios é atendido:
- O tamanho do lote atinge o tamanho máximo (100.000 linhas ou 32MB por 1GB de memória de réplica)
- O lote permanece aberto pelo tempo máximo permitido (5 segundos)
A latência (definida como o tempo entre o envio da mensagem do Kinesis para o stream e o momento em que ela fica disponível no ClickHouse) dependerá de vários fatores (por exemplo, a latência do Kinesis, a latência de rede e o tamanho/formato da mensagem). O envio em lotes descrito na seção acima também afeta a latência. Recomendamos sempre testar seu caso de uso específico para entender a latência esperada.
Se você tiver requisitos específicos de baixa latência, entre em contato conosco.
Recomendamos fortemente limitar o número de shards ativas simultaneamente de acordo com seus requisitos de throughput. Para um stream do Kinesis “On Demand”, a AWS atribui automaticamente um número correspondente de shards com base no throughput,
mas, para streams “Provisioned”, provisionar shards demais pode causar latência, como descrito abaixo, além de aumentar os custos, porque a cobrança do Kinesis para esses streams é feita “por shard”.
Se a sua aplicação produtora gravar continuamente em um grande número de shards ativas, isso poderá causar latência caso o pipe não esteja dimensionado adequadamente para processar essas shards com eficiência. Com base nos limites de throughput do Kinesis,
o ClickPipes atribui um número específico de “workers” por réplica para ler dados das shards. Por exemplo, no menor tamanho, uma réplica do ClickPipes terá 4 dessas threads de worker. Se o produtor estiver gravando
em mais de 4 shards ao mesmo tempo, os dados das shards “extras” não serão processados até que uma thread de worker fique disponível. Em particular, se o pipe estiver usando “enhanced fanout”, cada thread de worker ficará inscrita em uma
única shard por 5 minutos e não ficará disponível para ler nenhuma outra shard durante esse período. Isso pode causar “picos” de latência em múltiplos de 5 minutos.
O ClickPipes para Kinesis foi projetado para escalar tanto horizontal quanto verticalmente. Por padrão, criamos um grupo de consumidores com um único consumidor. Isso pode ser configurado durante a criação do ClickPipe ou, a qualquer momento, em Configurações -> Configurações avançadas -> Escalonamento.
O ClickPipes oferece alta disponibilidade com uma arquitetura distribuída entre zonas de disponibilidade.
Isso exige o escalonamento para pelo menos dois consumidores.
Independentemente do número de consumidores em execução, a tolerância a falhas é garantida por design.
Se um consumidor ou a infraestrutura subjacente falhar,
o ClickPipe reiniciará automaticamente o consumidor e continuará processando as mensagens.
Para acessar streams do Amazon Kinesis, você pode usar credenciais do IAM ou uma função do IAM. Para mais detalhes sobre como configurar uma função do IAM, você pode consultar este guia para obter informações sobre como configurar uma função que funcione com o ClickHouse Cloud