Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

iTechGuides is reader-supported. When you buy through links on our site, we may earn an affiliate commission. As an Amazon Associate I earn from qualifying purchases. Learn more

En Java, usa KafkaConsumer.subscribe(Pattern) para que el consumidor se suscriba a los tópicos cuyos nombres coincidan con una expresión regular. Por ejemplo, ^orders..* coincide con nombres que comienzan por orders.. Kafka revisa las coincidencias periódicamente: un tópico nuevo puede incorporarse después de actualizarse los metadatos y coordinarse el grupo, no necesariamente en el instante en que se crea.

Suscripción dinámica con subscribe(Pattern)

La API de consumidor de Apache Kafka permite suscribirse a todos los tópicos que coincidan con un patrón y recibir las particiones que el grupo asigne. En la API clásica de Java, compila la expresión con java.util.regex.Pattern y pásala a subscribe. La documentación de KafkaConsumer 4.1.2 describe esta opción como una suscripción a los tópicos que coinciden con el patrón para obtener particiones asignadas dinámicamente.

import java.time.Duration;
import java.util.regex.Pattern;
import org.apache.kafka.clients.consumer.KafkaConsumer;

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Pattern.compile("^orders\..*"));

while (running) {
    var records = consumer.poll(Duration.ofMillis(100));
    // Procesa records y gestiona commits según las necesidades de la aplicación.
}

En esta expresión, ^ marca el inicio del nombre, orders es el prefijo literal, . representa el punto literal y .* acepta cualquier secuencia posterior. Por tanto, nombres como orders.us y orders.eu.priority coinciden; order-events no. Es un ejemplo ilustrativo, no una prueba contra un clúster.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

props debe incluir, entre otras propiedades, la conexión al clúster y un group.id apropiado. El fragmento no muestra cierre del consumidor, gestión de errores ni una estrategia de commits: esos aspectos dependen de la aplicación.

Cómo llegan los tópicos nuevos y por qué puede haber rebalanceos

La suscripción por patrón no se evalúa necesariamente en cada mensaje. Kafka compara periódicamente el patrón con los tópicos existentes; metadata.max.age.ms influye en la frecuencia con que se refrescan los metadatos y se revisan esos cambios. Por ello, no hay que tratar el descubrimiento de un tópico recién creado como instantáneo.

Si cambia el conjunto de tópicos coincidentes, sus particiones o los miembros del grupo, Kafka puede rebalancear las asignaciones. El grupo, identificado por group.id, distribuye sus particiones entre los consumidores que lo integran. La aplicación debe continuar llamando a poll(Duration): la coordinación y los rebalances del grupo ocurren durante una llamada activa a poll, según la documentación de KafkaConsumer 4.1.2.

Elige la API según cómo cambia tu conjunto de tópicos

Necesidad API Consideración
Conjunto conocido y cerrado subscribe(Collection<String>) La lista es explícita y una nueva llamada reemplaza la suscripción anterior; no añade nombres a la lista previa.
Tópicos dinámicos seleccionados por nombre subscribe(Pattern) Kafka revisa periódicamente las coincidencias; los cambios pueden provocar rebalanceos.
Regex con el protocolo de consumidor nuevo en Kafka 4.x subscribe(SubscriptionPattern) Requiere group.protocol=consumer; el patrón debe ser compatible con RE2/J y se evalúa en el servidor.
Asignación manual de particiones assign(Collection<TopicPartition>) Omite la gestión del grupo y no se combina con la suscripción dinámica.

La diferencia de protocolo importa: Kafka documenta por separado la API clásica subscribe(Pattern) y SubscriptionPattern. El requisito de compatibilidad con RE2/J corresponde a esta última vía, no debe atribuirse automáticamente a la API clásica. Consulta el protocolo de rebalanceo de consumidor de Kafka 4.3 para la vía nueva.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

Gestión de offsets durante un rebalanceo

Si tu aplicación administra offsets o estado de particiones por su cuenta, usa el overload subscribe(Pattern, ConsumerRebalanceListener) para recibir callbacks cuando se asignen o revoquen particiones. El listener permite, por ejemplo, confirmar los offsets pendientes antes de que termine un rebalanceo. La gestión automática de offsets puede requerir otro enfoque; el listener es especialmente pertinente cuando el código controla esos commits.

Errores que conviene evitar

  • Esperar descubrimiento inmediato: el patrón se compara periódicamente con los tópicos existentes; la actualización de metadatos y la coordinación del grupo afectan el momento de incorporación.
  • Esperar procesamiento sin poll: la suscripción no sustituye al loop de consumo ni a las llamadas activas a poll(Duration).
  • Suponer que una lista se amplía: subscribe(Collection<String>) reemplaza la suscripción anterior.
  • Mezclar suscripción y asignación manual: assign(...) establece manualmente la asignación actual y omite la gestión del grupo; no es un complemento a subscribe(...). Sigue las reglas de la API para cancelar la suscripción al cambiar de modalidad.
  • Aplicar RE2/J a toda expresión de Java: esa condición corresponde a SubscriptionPattern con el protocolo de consumidor nuevo, no automáticamente a subscribe(Pattern).

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.