En el ecosistema de Apache Flume, las fuentes (sources) actúan como el punto de entrada para los eventos. Existen dos interfaces principales para implementar fuentes personalizadas: PollableSource y EventDrivenSource. La principal diferencia radica en su modelo de ejecución:
- PollableSource: Utiliza un modelo de "tirón" (pull). Un hilo de trabajo invoca periódicamente el método
process()para verificar si hay nuevos datos disponibles. - EventDrivenSource: Utiliza un modelo de "empuje" (push). La fuente espera a que ocurra un evento externo (como una conexión HTTP o un mensaje de cola) para disparar el procesamiento.
Para integrar Kafka con Flume utilizando un enfoque activo, es común implementar PollableSource. A continuación, se presenta una implementación robusta que consume mensajes JSON de Kafka, los transforma en eventos de Flume y los envía al canal. Esta clase también implementa Configurable para permitir la inyección de parámetros desde el archivo de configuración.
Código de la Fuente Customizada (Java)
import org.apache.flume.Context;
import org.apache.flume.Event;
import org.apache.flume.EventDeliveryException;
import org.apache.flume.PollableSource;
import org.apache.flume.conf.Configurable;
import org.apache.flume.event.EventBuilder;
import org.apache.flume.source.AbstractSource;
import com.alibaba.fastjson.JSONObject;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.nio.charset.StandardCharsets;
import java.text.SimpleDateFormat;
import java.util.*;
/**
* Fuente personalizada para ingesta de alertas EVM desde Kafka.
* Implementa PollableSource para consumo activo mediante polling.
*/
public class EvAlarmKafkaSource extends AbstractSource implements Configurable, PollableSource {
private static final String TOPIC_KEY = "topic";
private static final String GROUP_ID_KEY = "groupId";
private static final String BROKER_LIST_KEY = "brokerList";
private String topicName;
private String consumerGroupId;
private String bootstrapServers;
private KafkaConsumer<String, String> kafkaClient;
private SimpleDateFormat dateFormatter;
@Override
public void configure(Context context) {
// Recuperación de parámetros desde flume.conf
this.topicName = context.getString(TOPIC_KEY);
this.consumerGroupId = context.getString(GROUP_ID_KEY);
this.bootstrapServers = context.getString(BROKER_LIST_KEY);
// Inicialización del formato de fecha para metadatos
this.dateFormatter = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
}
/**
* Método llamado repetidamente por el hilo de la fuente para procesar eventos.
*/
@Override
public Status process() throws EventDeliveryException {
if (kafkaClient == null) {
return Status.BACKOFF;
}
try {
// Polling con timeout corto para evitar bloqueos largos
ConsumerRecords<String, String> records = kafkaClient.poll(100);
for (ConsumerRecord<String, String> record : records) {
try {
String rawPayload = record.value();
// Parseo del JSON recibido de Kafka
JSONObject jsonNode = JSONObject.parseObject(rawPayload);
// Extracción de campos específicos de la alerta
String messageId = jsonNode.getString("msgid");
String vehicleVin = jsonNode.getString("vin");
long creationTimeSec = jsonNode.getLongValue("tm_c") != 0 ?
jsonNode.getLongValue("tm_c") : System.currentTimeMillis()/1000;
long updateTimeSec = jsonNode.getLongValue("tm_u") != 0 ?
jsonNode.getLongValue("tm_u") : System.currentTimeMillis()/1000;
String typeFlag = jsonNode.getString("t_flg");
String hazardLevel = jsonNode.getString("ha_level");
String caFlag = jsonNode.getString("ca_flg");
String selfValue = jsonNode.getString("sef_va");
String dmfValue = jsonNode.getString("dmf_va");
String egfValue = jsonNode.getString("egf_va");
String ofValue = jsonNode.getString("of_va");
String longitude = jsonNode.getString("lng");
String latitude = jsonNode.getString("lat");
// Construcción del cuerpo del evento (delimitado por #)
StringBuilder payloadBuilder = new StringBuilder();
payloadBuilder.append(messageId).append("#")
.append(vehicleVin).append("#")
.append(dateFormatter.format(new Date(creationTimeSec * 1000))).append("#")
.append(dateFormatter.format(new Date(updateTimeSec * 1000))).append("#")
.append(typeFlag).append("#")
.append(hazardLevel).append("#")
.append(caFlag).append("#")
.append(selfValue).append("#")
.append(dmfValue).append("#")
.append(egfValue).append("#")
.append(ofValue).append("#")
.append(longitude).append("#")
.append(latitude);
// Preparación de encabezados del evento Flume
Map<String, String> headers = new HashMap<>();
headers.put("timestamp", String.valueOf(creationTimeSec * 1000));
headers.put("source_type", "kafka_ev_alarm");
// Creación y envío del evento al ChannelProcessor
Event flumeEvent = EventBuilder.withBody(payloadBuilder.toString(), StandardCharsets.UTF_8, headers);
getChannelProcessor().processEvent(flumeEvent);
} catch (Exception e) {
// Registro de error sin detener el flujo completo
System.err.println("Error procesando registro individual: " + e.getMessage());
}
}
return Status.READY;
} catch (Exception e) {
throw new EventDeliveryException("Fallo crítico durante el polling de Kafka", e);
}
}
@Override
public synchronized void start() {
super.start();
try {
Properties props = new Properties();
props.put("bootstrap.servers", bootstrapServers);
props.put("group.id", consumerGroupId);
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "1000");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
kafkaClient = new KafkaConsumer<>(props);
kafkaClient.subscribe(Collections.singletonList(topicName));
} catch (Exception e) {
throw new RuntimeException("No se pudo inicializar el consumidor de Kafka", e);
}
}
@Override
public synchronized void stop() {
if (kafkaClient != null) {
kafkaClient.close();
kafkaClient = null;
}
super.stop();
}
@Override
public long getBackOffSleepIncrement() {
return 1000;
}
@Override
public long getMaxBackOffSleepInterval() {
return 5000;
}
}
Configuración del Agente Flume
Una vez compilada la fuente, debe registrarse en el archivo flume.conf. En este ejemplo, los datos procesados se escriben en HDFS. Es crucial ajustar los parámetros de rotación de archivso (rollSize, rollCount, rollInterval) según el volumen de datos esperado.
# Definición de componentes
agent.sources = ev_source
agent.channels = mem_channel
agent.sinks = hdfs_sink
# Configuración de la Fuente Customizada
agent.sources.ev_source.type = com.example.flume.EvAlarmKafkaSource
agent.sources.ev_source.channels = mem_channel
# Parámetros pasados a la interfaz Configurable
agent.sources.ev_source.topic = evm-alerts-topic
agent.sources.ev_source.groupId = flume-consumer-group-01
agent.sources.ev_source.brokerList = kafka-broker-1:9092,kafka-broker-2:9092
# Configuración del Canal (Memory Channel para baja latencia)
agent.channels.mem_channel.type = memory
agent.channels.mem_channel.capacity = 200000
agent.channels.mem_channel.transactionCapacity = 20000
agent.channels.mem_channel.byteCapacityBufferPercentage = 20
# Configuración del Sink (HDFS)
agent.sinks.hdfs_sink.type = hdfs
agent.sinks.hdfs_sink.channel = mem_channel
agent.sinks.hdfs_sink.hdfs.path = hdfs://nameservice1/flume/evm/alarm/%Y%m%d
agent.sinks.hdfs_sink.hdfs.fileType = DataStream
agent.sinks.hdfs_sink.hdfs.writeFormat = Text
agent.sinks.hdfs_sink.hdfs.fileSuffix = .log
# Estrategia de Rotación de Archivos
# Rota cada 60 minutos o cuando alcanza 256MB, independientemente del conteo de eventos
agent.sinks.hdfs_sink.hdfs.rollInterval = 3600
agent.sinks.hdfs_sink.hdfs.rollSize = 268435456
agent.sinks.hdfs_sink.hdfs.rollCount = 0
agent.sinks.hdfs_sink.hdfs.batchSize = 100
# Tiempo de inactividad antes de cerrar el archivo temporal (0 deshabilita esta lógica específica)
agent.sinks.hdfs_sink.hdfs.idleTimeout = 0
# Replicación y Timeouts
agent.sinks.hdfs_sink.hdfs.minBlockReplicas = 1
agent.sinks.hdfs_sink.hdfs.callTimeout = 120000
agent.sinks.hdfs_sink.hdfs.useLocalTimeStamp = true
Dsepliegue y Ejecución
- Empaquetado: Compile la clase Java y empaquétela en un archivo JAR. Asegúrese de incluir todas las dependencias (Kafka Clients, FastJSON) o utilice un plugin Fat Jar.
- Instalación: Copie el JAR generado al directorio
lib/dentro de la instalación de Flume en todos los nodos del clúster donde correrá el agente. - Ejecución: Inicie el agente especificando la ruta de configuración y el nombre del agente definido en el archivo conf.