AWS Glue Streaming
O AWS Glue Streaming, um componente do AWS Glue, possibilita que você lide com dados de streaming quase em tempo real com eficiência, capacitando-o a realizar tarefas cruciais, como a ingestão de dados, o processamento e o machine learning. Ao usar a estrutura do Apache Spark Streaming, o AWS Glue Streaming fornece um serviço sem servidor que pode lidar com dados de streaming em escala. O AWS Glue disponibiliza diversas otimizações além do Apache Spark, como infraestrutura sem servidor, ajuste de escala automático, desenvolvimento visual de trabalhos, cadernos instantâneos para trabalhos de streaming e outros aprimoramentos de performance.
Casos de uso para o streaming
Alguns casos de uso comuns para o AWS Glue Streaming incluem:
Processamento de dados quase em tempo real: o AWS Glue Streaming permite que as organizações processem dados de streaming quase em tempo real, permitindo-lhes obter insights e tomar decisões oportunas com base nas informações mais recentes.
Detecção de fraudes: é possível utilizar o AWS Glue Streaming para realizar análises em tempo real de dados de streaming, tornando-o valioso para a detecção de atividades fraudulentas, como a fraude de um cartão de crédito, a invasão de rede ou as fraudes on-line. Ao processar e analisar continuamente os dados de entrada, você pode identificar rapidamente padrões ou anomalias suspeitas.
Analytics de mídia social: o AWS Glue Streaming pode processar dados de mídia social em tempo real, como tweets, publicações ou comentários, possibilitando que as organizações monitorem tendências, analisem sentimentos e gerenciem a reputação da marca em tempo real.
Analytics da Internet das Coisas (IoT): o AWS Glue Streaming é adequado para analisar e lidar com fluxos de dados de alta velocidade gerados por dispositivos de IoT, sensores e máquinas conectadas. Ele permite o monitoramento em tempo real, a detecção de anomalias, a manutenção preditiva e outros casos de uso de analytics da IoT.
Análise de fluxo de cliques: o AWS Glue Streaming pode processar e analisar dados de fluxo de cliques em tempo real de sites ou de aplicações móveis. Isso possibilita que as empresas obtenham insights sobre o comportamento do usuário, personalizem as experiências do usuário e otimizem campanhas de marketing com base em dados de fluxo de cliques em tempo real.
Monitoramento e análise de log: o AWS Glue Streaming pode processar e analisar continuamente os dados de log de servidores, aplicações ou dispositivos de rede em tempo real. Isso ajuda a detectar anomalias, solucionar problemas e monitorar a integridade e a performance do sistema.
Sistemas de recomendação: o AWS Glue Streaming pode processar dados de atividades do usuário em tempo real e atualizar os modelos de recomendação de forma dinâmica. Isso permite recomendações personalizadas e em tempo real com base no comportamento e nas preferências do usuário.
Esses são alguns exemplos da diversidade de casos de uso em que o AWS Glue Streaming pode ser aplicado. Sua integração com o ecossistema e os serviços gerenciados da AWS o torna uma escolha conveniente para o processamento e para a analytics de fluxo em tempo real na nuvem.
Quais são os benefícios do uso do AWS Glue Streaming?
Os benefícios do uso do AWS Glue Streaming são os seguintes:
Tecnologia sem servidor: o AWS Glue Streaming tem tecnologia sem servidor, o que elimina a necessidade de gerenciamento da infraestrutura. Isso reduz a sobrecarga operacional e permite que os usuários se concentrem nas tarefas de processamento e de analytics de dados, em vez de no gerenciamento da infraestrutura.
Ajuste de escala automático: o AWS Glue Streaming fornece recursos de ajuste de escala automático, ajustando dinamicamente a capacidade de processamento com base na workload. Ele aumenta ou reduz a escala horizontalmente de forma automática para lidar com as flutuações no volume de dados, garantindo uma performance e uma utilização de recursos ideais.
Desenvolvimento visual: o desenvolvimento do trabalho de streaming pode ser complexo. O AWS Glue Streaming aborda esse desafio ao disponibilizar o AWS Glue Studio, uma ferramenta de criação visual. O AWS Glue Studio simplifica o processo de criação de fluxos de trabalho de streaming e possibilita que os desenvolvedores projetem e gerenciem aplicações de streaming visualmente, reduzindo a curva de aprendizado e aumentando a produtividade.
Econômico: como um serviço sem servidor, o AWS Glue Streaming oferece uma relação custo-benefício vantajosa ao eliminar a necessidade de provisionamento e manutenção de infraestrutura. Os usuários são cobrados com base nos recursos consumidos durante a execução de trabalhos de streaming, permitindo a otimização de custos e a escalabilidade com base no uso real.
Tratamento de workloads complexas: o AWS Glue Streaming foi projetado para lidar com workloads de streaming complexas. Ele pode processar e analisar grandes volumes de dados em tempo real, oferecer suporte a transformações avançadas e se integrar a outros serviços da AWS, possibilitando pipelines de dados de streaming e fluxos de trabalho de analytics sofisticados.
Sem aprisionamento tecnológico: o AWS Glue Streaming oferece flexibilidade e evita o aprisionamento tecnológico do fornecedor. Os usuários podem aproveitar o AWS Glue Streaming como parte do ecossistema mais amplo da AWS, integrando-o sem complicações a outros serviços da AWS. Isso permite a fácil integração com fontes de dados, aplicações e serviços existentes, sem a necessidade de estar vinculado a uma tecnologia ou plataforma específica.
Quando devo usar o AWS Glue Streaming?
Existem muitas opções quando se trata de casos de uso de streaming. Recomendamos o uso do streaming do AWS Glue nos cenários apresentados a seguir.
Se você já usa o AWS Glue ou o Spark para realizar o processamento em lote, o AWS Glue Streaming é a escolha ideal para você. Ele fornece uma transição sem complicação para o desenvolvimento de trabalhos de streaming sem a necessidade de aprender uma nova linguagem ou estrutura. Aproveitando o conhecimento e a infraestrutura existentes, o AWS Glue Streaming simplifica o processo de desenvolvimento de trabalhos e permite ampliar os recursos de processamento de dados com facilidade para cenários de streaming em tempo real.
Se você precisar de um serviço ou de um produto unificado para lidar com workloads em lote, de streaming e orientadas a eventos, o AWS Glue Streaming é a solução para você. Com o AWS Glue Streaming, é possível consolidar suas necessidades de processamento de dados em uma única estrutura, eliminando a complexidade do gerenciamento de vários sistemas. Isso possibilita o desenvolvimento e a manutenção eficientes de diversos fluxos de trabalho de dados, ao mesmo tempo em que garante consistência e compatibilidade entre diferentes tipos de workloads.
O AWS Glue Streaming é adequado para cenários que envolvem volumes de dados de streaming extremamente grandes e transformações complexas, como junções entre fluxos ou bancos de dados relacionais. Ele pode processar e analisar fluxos massivos de dados com eficiência, possibilitando que você lide com workloads complexas com facilidade. Quer se trate de ingestão de dados em alta velocidade ou de manipulações complexas de dados, a escalabilidade e os recursos avançados de processamento do AWS Glue Streaming garantem performance ideal e resultados precisos.
Se você preferir uma abordagem visual para desenvolver trabalhos de streaming, o AWS Glue oferece o AWS Glue Studio, com o qual você pode projetar e gerenciar visualmente as aplicações de streaming, simplificando o processo de desenvolvimento. Essa interface intuitiva possibilita que os desenvolvedores criem, configurem e monitorem fluxos de trabalho de streaming usando uma interface visual, o que reduz a curva de aprendizado e aumenta a produtividade.
O AWS Glue Streaming é uma excelente opção para casos de uso quase em tempo real, nos quais existem SLAs (Acordos de Serviço) rigorosos superiores a dez segundos.
Se você estiver desenvolvendo um data lake transacional usando o Apache Iceberg, o Apache Hudi ou o Delta Lake, o AWS Glue Streaming fornecerá suporte nativo para esses formatos de tabela aberta. Essa integração sem complicações possibilita o processamento de dados de streaming diretamente desses data lakes transacionais, garantindo integridade, compatibilidade e consistência de dados.
Ao precisar ingerir dados de streaming para uma variedade de destinos de dados, o AWS Glue Streaming disponibiliza destinos nativos para uma variedade de destinos de dados, como o Amazon Redshift, o Amazon RDS, o Amazon Aurora, o Oracle, o SQL Server e outros destinos.
Fontes de dados compatíveis
AWS GlueO Streaming oferece suporte às seguintes fontes de dados:
Amazon Kinesis
Amazon MSK (Managed Streaming for Apache Kafka)
Apache Kafka autogerenciado
Destinos de dados com suporte
AWS GlueO Streaming oferece suporte para uma variedade de destinos de dados, como:
Destinos de dados compatíveis com o Catálogo de Dados do AWS Glue
Amazon S3
Amazon Redshift
MySQL
PostgreSQL
Oracle
Microsoft SQL Server
Snowflake
Qualquer banco de dados que possa ser conectado usando JDBC
Apache Iceberg, Delta e Apache Hudi
AWS GlueConectores do Marketplace
Habilitar o modo de tempo real para trabalhos de streaming
O modo de tempo real (RTM) é um novo modelo de execução do Spark Structured Streaming disponível no AWS Glue versão 6.0. O RTM reduz a latência de ponta a ponta de segundos ou minutos para menos de um segundo. O modo de tempo real se aplica somente aos trabalhos do Spark Structured Streaming. Ele não se aplica ao Spark Streaming antigo (DStreams) ou a outros tipos de trabalho.
O RTM usa Trigger.RealTime. As tarefas são executadas continuamente em uma janela de lote (padrão de cinco minutos) e processam os registros à medida que chegam, em vez de acumular dados entre os intervalos. Isso difere do modelo de microlote padrão, em que forEachBatch/Trigger.ProcessingTime pesquisa, processa, confirma e reinicia tarefas a cada intervalo.
Importante
O RTM requer sua aceitação explícita por meio de um argumento de trabalho. Se não houver slots de tarefas suficientes para todas as partições de origem, o RTM descartará silenciosamente as partições não atribuídas. Você deve provisionar operadores suficientes para todas as partições Kafka.
Pré-requisitos
Antes de habilitar o modo de tempo real, confirme se o trabalho atende aos seguintes requisitos:
-
AWS Glue versão 6.0
-
O trabalho deve usar o Spark Structured Streaming. O modo de tempo real não se aplica ao Spark Streaming antigo (DStreams) ou outros tipos de trabalho.
-
O tipo de trabalho deve ser Spark Streaming (comando
gluestreaming) -
A linguagem do trabalho deve ser Scala (
--job-language scala). A compatibilidade com o RTM do PySpark não estará disponível até o Spark 4.2. -
Fonte Kafka apenas. O Amazon Kinesis não é compatível com RTM no AWS Glue versão 6.0.
-
Operações sem estado (selecionar, filtrar, projetar, mapear) apenas. Operações com estado, como agregações, uniões, desduplicação e operações com janelas, não são compatíveis.
-
O modo de saída deve ser Atualização. O modo Acréscimo não é compatível com o RTM.
-
O ajuste de escala automático não é compatível com o modo de tempo real. Não habilite o ajuste de escala automático para trabalhos do RTM. Configure um número fixo de operadores suficiente para todas as partições Kafka no tópico de origem.
Quando usar o modo de tempo real
O modo de tempo real foi projetado para uma classe específica de workloads de streaming. Considere o uso do modo em tempo real quando:
-
Você precisar de latência de ponta a ponta abaixo de um segundo, e a latência de microlotes (1 a 2 segundos ou mais) for muito alta para o caso de uso.
-
O pipeline realizar transformações sem estado, como filtrar, projetar, enriquecer ou rotear registros de Kafka para Kafka ou outro coletor.
-
Você tiver um número fixo e previsível de partições Kafka e puder provisionar operadores adequadamente.
-
Os trabalhos forem escritos em Scala.
Continue usando o modo de microlote quando:
-
Você precisar de operações com estado, como agregações, uniões, desduplicação ou cálculos em janelas.
-
Você usar o Amazon Kinesis como uma origem.
-
Você escrever trabalhos em PySpark.
-
Você precisar de ajuste de escala automático para lidar com volumes de dados variáveis.
-
Você usar a API de streaming
forEachBatchou GlueContext. -
Uma latência de segundos for aceitável para o caso de uso.
Como o modo de tempo real funciona
A seguir, descrevemos a diferença entre o modo de microlote e o modo de tempo real:
- Modo de microlote
-
Cada intervalo inicia as tarefas, lê os dados acumulados, processa os dados, confirma o ponto de verificação, encerra as tarefas e repete. A latência mínima é aproximadamente de 1 a 2 segundos.
- Modo de tempo real
-
As tarefas são iniciadas uma vez e são executadas pela duração de
batchDurationMs(padrão 5 minutos). As tarefas processam os registros à medida que chegam, com latência abaixo de um segundo. No fim do prazo, as tarefas são interrompidas de modo cooperativo. O driver confirma o ponto de verificação e o próximo lote reinicia as tarefas.
Ambos os modos usam o mesmo formato de ponto de verificação e mecanismo de recuperação. A principal diferença é a vida útil da tarefa. O modo de microlote encerra e reinicia as tarefas a cada intervalo. O modo de tempo real mantém as tarefas em execução contínua dentro de uma janela de lote maior.
Importante
Se não houver slots de suficientes para processar todas as partições de origem, o RTM descartará silenciosamente as partições não atribuídas. Certifique-se de provisionar operadores suficientes para cobrir todas as partições.
Para habilitar o modo de tempo real
Você habilita o modo em tempo real definindo o argumento do trabalho --enable-real-time-mode como true. É possível definir esse argumento no console do AWS Glue ou por meio da API.
Para habilitar o modo de tempo real (console)
-
Abra o console do AWS Glue
e abra o trabalho de streaming. -
Escolha a guia Job details (Detalhes do trabalho).
-
Em Versão do Glue, selecione Glue 6.0. Em Tipo, escolha Spark Streaming.
-
Role para baixo até a seção Parâmetros do trabalho.
-
Escolha Adicionar novos parâmetros.
-
Em Chave, digite
--enable-real-time-mode. Em Valor, insiratrue. -
Escolha Salvar.
nota
Os traços iniciais são obrigatórios. Parâmetros do trabalho é a visualização do console de DefaultArguments.
Para habilitar o modo de tempo real (API)
O sinalizador --enable-real-time-mode é armazenado no mapa DefaultArguments da definição do trabalho. É possível configurá-lo ao criar ou atualizar um trabalho.
Para criar um novo trabalho (AWS CLI)
Execute o seguinte comando:
aws glue create-job \ --name my-rtm-job \ --role arn:aws:iam::123456789012:role/MyGlueRole \ --glue-version 6.0 \ --worker-type G.1X --number-of-workers 4 \ --command '{"Name":"gluestreaming","ScriptLocation":"s3://my-bucket/scripts/rtm-job.scala"}' \ --default-arguments '{ "--enable-real-time-mode": "true", "--job-language": "scala", "--class": "GlueApp", "--TempDir": "s3://my-bucket/tmp/" }' \ --region us-east-2
Para criar um novo trabalho (boto3)
Use o seguinte código:
import boto3 glue = boto3.client("glue", region_name="us-east-2") glue.create_job( Name="my-rtm-job", Role="arn:aws:iam::123456789012:role/MyGlueRole", GlueVersion="6.0", WorkerType="G.1X", NumberOfWorkers=4, Command={ "Name": "gluestreaming", "ScriptLocation": "s3://my-bucket/scripts/rtm-job.scala", }, DefaultArguments={ "--enable-real-time-mode": "true", "--job-language": "scala", "--class": "GlueApp", "--TempDir": "s3://my-bucket/tmp/", }, )
Para atualizar um trabalho existente (AWS CLI)
Execute o seguinte comando:
aws glue update-job \ --job-name my-existing-job \ --job-update '{ "GlueVersion": "6.0", "DefaultArguments": { "--enable-real-time-mode": "true", "--job-language": "scala" } }'
Escrever o script de streaming
O argumento do trabalho declara a intenção de usar o modo de tempo real. O script seleciona o acionador.
O seguinte exemplo em Scala mostra uma consulta de streaming que usa Trigger.RealTime:
import org.apache.spark.sql.streaming.Trigger val query = df.writeStream .format("kafka") .outputMode("update") .trigger(Trigger.RealTime(60000L)) // checkpoint interval in milliseconds .start() query.awaitTermination()
O Trigger.RealTime usa um intervalo de ponto de verificação em milissegundos. O modo de saída da atualização é obrigatório. O modo de adição gera OUTPUT_MODE_NOT_SUPPORTED.
Você pode misturar modos em um script, desde que o sinalizador seja definido:
dfA.writeStream.outputMode("update").trigger(Trigger.RealTime(60000L)).start() dfB.writeStream.outputMode("append").trigger(Trigger.ProcessingTime("30 seconds")).start()
Comportamento quando não há sinalizador
A seguir é descrito como o trabalho se comporta quando o sinalizador --enable-real-time-mode não é definido:
-
Um trabalho que inicia uma consulta em tempo real sem o sinalizador
--enable-real-time-modefalha no início da consulta. A mensagem de falha orienta você a adicionar o argumento. -
Trabalhos somente em microlotes nunca são afetados pela ausência desse sinalizador.
-
Um trabalho que define o sinalizador, mas usa apenas consultas em microlotes, também não é afetado.
Considerações e limitações
Considere o seguinte quando usar o modo de tempo real:
- Partições descartadas
-
Se não houver slots de suficientes para processar todas as partições de origem, partições não atribuídas não serão processadas. Provisione operadores suficientes para todas as partições Kafka.
- Não há ajuste de escala automático
-
Não habilite ajuste de escala automático para trabalhos no modo de tempo real. O ajuste de escala automático não é compatível com o RTM e introduz uma latência que anula os benefícios da baixa latência. Provisione um número fixo de operadores igual ou maior que o número de partições Kafka no tópico de origem.
- Kafka apenas
-
A origem no Amazon Kinesis não é compatível com o RTM no AWS Glue versão 6.0.
- Scala apenas
-
O PySpark não é compatível com RTM até o Spark 4.2.
- Sem estado apenas
-
Agregações, uniões, desduplicação e operações com janelas e
transformWithStatenão são compatíveis. - forEachBatch incompatível
-
O RTM não usa o modelo
forEachBatch. UsewriteStreamcomTrigger.RealTimediretamente. - Recuperação de ponto de verificação
-
Na reinicialização do trabalho, o RTM faz a recuperação do último ponto de verificação. Os pontos de verificação ocorrem a cada
batchDurationMs. Na pior das hipóteses, o reprocessamento tem a duração de uma janela de lote (semântica pelo menos uma vez).