Otimizar e monitorar o desempenho da consulta no modo em tempo real

Esta página aborda o ajuste de computação, técnicas para reduzir a latência de ponta a ponta e abordagens para medir o desempenho da consulta no modo em tempo real.

Otimização de computação

Ao configurar sua computação, considere o seguinte:

  • Ao contrário do modo de microlote, as tarefas em tempo real podem permanecer ociosas enquanto aguardam dados, portanto, o dimensionamento correto é essencial para evitar o desperdício de recursos.
  • Objetive um nível de utilização de cluster desejado, como 50%, ajustando:
    • maxPartitions (para Kafka)
    • spark.sql.shuffle.partitions (para estágios de embaralhamento)
  • O Databricks recomenda a configuração maxPartitions para que cada tarefa trate várias partições kafka para reduzir a sobrecarga.
  • Ajuste os slots de tarefa por trabalhador para corresponder à carga de trabalho para tarefas simples de estágio único.
  • No caso de trabalhos com muito embaralhamento, faça experiências para encontrar o número mínimo de partições de embaralhamento que evitem listas de pendências e ajuste a partir daí. O computador não agendará a tarefa se não tiver slots suficientes.

Observação

A partir do Databricks Runtime 16.4 LTS e superior, todos os pipelines em tempo real usam o ponto de verificação v2 para permitir alternâncias entre os modos em tempo real e microlote.

Otimização de latência

O modo de streaming estruturado em tempo real tem técnicas opcionais para reduzir a latência de ponta a ponta. Nenhum dos dois está habilitado por padrão. Você deve ativá-los cada um separadamente.

  • Acompanhamento de progresso assíncrono: move gravações para o deslocamento e os logs de commits em um thread assíncrono, reduzindo o tempo entre lotes durante as consultas sem estado.
  • Ponto de verificação de estado assíncrono: começa a processar o próximo microlote assim que a computação for concluída, sem aguardar o ponto de verificação de estado, reduzindo a latência para consultas com estado.

Monitoramento e observabilidade

No modo em tempo real, as métricas tradicionais de duração do lote não refletem a latência real de ponta a ponta. Use as abordagens abaixo para medir a latência com precisão e identificar gargalos em suas consultas.

A latência de ponta a ponta é específica da carga de trabalho e, às vezes, só pode ser medida com precisão com a lógica de negócios. Por exemplo, se o carimbo de data/hora de origem for gerado no Kafka, você pode calcular a latência como a diferença entre o carimbo de data/hora de saída do Kafka e o carimbo de data/hora de origem.

Métricas integradas com StreamingQueryProgress

O evento StreamingQueryProgress é registrado automaticamente nos logs do driver e acessível por meio da função de retorno de chamada StreamingQueryListeneronQueryProgress(). Isso permite que você reaja a eventos de progresso programaticamente, por exemplo, se quiser publicar métricas em um sistema de monitoramento externo. QueryProgressEvent.json() ou toString() inclua estas métricas de modo em tempo real:

  1. Latência de processamento (processingLatencyMs). O tempo transcorrido entre quando a consulta no modo em tempo real lê um registro e quando grava esse registro no próximo estágio ou downstream. O sistema relata essa métrica por tarefa.
  2. Latência de fila de origem (sourceQueuingLatencyMs). A quantidade de tempo decorrido entre quando o sistema grava um registro em um barramento de mensagens—por exemplo, o tempo de anexação do log no Kafka—e quando a consulta no modo em tempo real lê o registro pela primeira vez. O sistema relata essa métrica por tarefa.
  3. Latência de ponta a ponta (e2eLatencyMs). O intervalo entre o momento em que o sistema grava o registro em um barramento de mensagens e quando a consulta no modo de tempo real grava o registro a jusante. O sistema agrega essa métrica por lote em todos os registros processados por todas as tarefas.

O evento de progresso JSON inclui essas métricas sob latencies. Por exemplo:

{
  "latencies": {
    "processingLatencyMs": {
      "P0": 0,
      "P50": 0,
      "P90": 0,
      "P95": 0,
      "P99": 0
    },
    "sourceQueuingLatencyMs": {
      "P0": 0,
      "P50": 1,
      "P90": 1,
      "P95": 2,
      "P99": 3
    },
    "e2eLatencyMs": {
      "P0": 0,
      "P50": 1,
      "P90": 1,
      "P95": 2,
      "P99": 4
    }
  }
}

Monitorar a utilização das tarefas

Use busyTimeFraction para determinar se a utilização das tarefas do Spark está limitando o rendimento. Valores próximos de 1 indicam que as tarefas estão totalmente utilizadas e a consulta pode precisar de mais cálculo. Valores mais baixos indicam que as tarefas passam mais tempo ociosas ou bloqueadas. A métrica varia de 0 a 1 e é reportada por estágio e tarefa.

Ative as métricas de depuração antes de iniciar a consulta de streaming:

spark.conf.set("spark.databricks.streaming.execution.enableDebugMetrics", "true")

Use um dos seguintes métodos para acessar métricas de tarefa:

Dados brutos

Inicie ou reinicie a consulta. Após o disparo ser concluído, abra Dados Brutos abaixo da célula de consulta e encontre _taskMetrics sob latencies. O exemplo a seguir mostra as métricas de tarefa:

{
  "latencies": {
    "_taskMetrics": {
      "stage_0_task_0": {
        "busyTimeFraction": 0.03
      }
    }
  }
}

Python

Após a conclusão de um gatilho, acesse _taskMetrics através de lastProgress:

import json

task_metrics = json.loads(query.lastProgress.json)["latencies"]["_taskMetrics"]
print(task_metrics)

Medida de latência personalizada com a API de Observação

A API Observe permite medir a latência embutida sem iniciar um trabalho separado. Se você tiver um timestamp de origem que se aproxime da hora de chegada dos dados de origem, poderá estimar a latência por lote registrando um timestamp antes do coletor e calculando a diferença. Os resultados aparecem em relatórios de progresso e estão disponíveis para ouvintes.

Python

from datetime import datetime

from pyspark.sql.functions import avg, col, lit, max, percentile_approx, udf, unix_millis
from pyspark.sql.types import TimestampType

@udf(returnType=TimestampType())
def current_timestamp():
  return datetime.now()

# Query before outputting
.withColumn("temp-timestamp", current_timestamp())
.withColumn(
  "latency",
  unix_millis(col("temp-timestamp")).cast("long") - unix_millis(col("timestamp")).cast("long"))
.observe(
  "observedLatency",
  avg(col("latency")).alias("avg"),
  max(col("latency")).alias("max"),
  percentile_approx(col("latency"), lit(0.99), lit(150)).alias("p99"),
  percentile_approx(col("latency"), lit(0.5), lit(150)).alias("p50"))
.drop(col("latency"))
.drop(col("temp-timestamp"))
# Output part of the query. For example, .WriteStream, etc.

Scala

import org.apache.spark.sql.functions.{avg, col, lit, max, percentile_approx, udf, unix_millis}

val currentTimestampUDF = udf(() => System.currentTimeMillis())

// Query before outputting
.withColumn("temp-timestamp", currentTimestampUDF())
.withColumn(
  "latency",
  col("temp-timestamp").cast("long") - unix_millis(col("timestamp")).cast("long"))
.observe(
  name = "observedLatency",
  avg(col("latency")).as("avg"),
  max(col("latency")).as("max"),
  percentile_approx(col("latency"), lit(0.99), lit(150)).as("p99"),
  percentile_approx(col("latency"), lit(0.5), lit(150)).as("p50"))
.drop(col("latency"))
.drop(col("temp-timestamp"))
// Output part of the query. For example, .WriteStream, etc.

Saída de exemplo:

"observedMetrics" : {
  "observedLatency" : {
    "avg" : 63.8369765176552,
    "max" : 219,
    "p99" : 154,
    "p50" : 49
  }
}

Recursos adicionais