Flink Time Characteristic e AutoWatermarkInterval
No Apache Flink, setAutoWatermarkInterval(interval)produz marcas d'água para operadores downstream para que avancem no horário do evento.
Se a marca d'água não foi alterada durante o intervalo especificado (nenhum evento chegou), o tempo de execução não emitirá nenhuma marca d'água? Por outro lado, se um novo evento chegar antes do próximo intervalo, uma nova marca d'água será emitida imediatamente ou será colocada na fila / aguardando até que o próximo intervalo setAutoWatermarkInterval seja alcançado.
Estou curioso para saber qual é a melhor configuração AutoWatermarkInterval (especialmente para fontes de alta taxa): quanto mais este valor for pequeno, maior será a defasagem entre o tempo de processamento e o tempo do evento, mas com a sobrecarga de mais uso de BW para enviar as marcas d'água . Isso é verdade, correto?
Por outro lado, se eu usar env.setStreamTimeCharacteristic (TimeCharacteristic.IngestionTime), o Flink runtime atribuirá automaticamente timestamps e watermarks (timestamps correspondem ao horário em que o evento entrou no pipeline de fluxo de dados do Flink, ou seja, o operador de origem), no entanto, mesmo com ingestionTime, podemos ainda definir um temporizador de tempo de processamento (na função processElement) como mostrado abaixo:
long timer = context.timestamp() + Timeout.
context.timerService().registerProcessingTimeTimer(timer);
onde context.timestamp () é o tempo de ingestão definido pelo Flink.
Obrigado.
Respostas
O autoWatermarkInterval afeta apenas os geradores de marca d'água que prestam atenção a ele. Eles também têm a oportunidade de gerar uma marca d'água em combinação com o processamento de eventos.
Para aqueles geradores de marca d'água que usam o autoWatermarkInterval (que é definitivamente o caso normal), eles estão coletando evidências de qual deve ser a próxima marca d'água como um efeito colateral da atribuição de carimbos de data / hora para cada evento. Quando um cronômetro dispara (com base no autoWatermarkInterval), o gerador de marca d'água é solicitado pelo tempo de execução do Flink para produzir a próxima marca d'água. A marca d'água não estava esperando em algum lugar, nem estava enfileirada, mas sim criada sob demanda, com base nas informações que foram armazenadas pelo designador de carimbo de data / hora - que normalmente é o carimbo de data / hora máximo visto até agora no fluxo.
Sim, marcas d'água mais frequentes significam mais sobrecarga para comunicá-las e processá-las, e menor latência. Você deve decidir como lidar com essa compensação de taxa de transferência / latência com base nos requisitos do seu aplicativo.
Você sempre pode usar temporizadores de tempo de processamento, independentemente do TimeCharacteristic. (A propósito, em um nível baixo, a única coisa que as marcas d'água fazem é acionar temporizadores de tempo de evento, sejam eles em funções de processo, janelas, etc.)