基于Java自定义Flume Source与Sink对接Kafka及Hadoop的技术咨询
Hey there! Since you’ve already got the config-based Flume-Kafka-HDFS pipeline running smoothly with your stack (Flume 1.6, Kafka 1.0.0, ZooKeeper 3.4.10), let’s break down how to build and deploy custom Java implementations for the Kafka Source and HDFS Sink—plus tackle the common issues you might be hitting in your Cloudera VM.
First, let’s build a source that pulls messages from Kafka using the modern Kafka 1.0.0 consumer API (since Flume 1.6’s native Kafka source uses outdated 0.8.x APIs that won’t play nicely with your setup).
1. Dependencies (Maven)
Make sure your pom.xml includes these compatible dependencies—note that we’re using the official kafka-clients jar, not the legacy Scala-based Kafka jars:
<dependencies> <!-- Flume Core --> <dependency> <groupId>org.apache.flume</groupId> <artifactId>flume-ng-core</artifactId> <version>1.6.0</version> <scope>provided</scope> <!-- Cloudera provides this --> </dependency> <!-- Kafka Clients (matches your 1.0.0 version) --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>1.0.0</version> </dependency> <!-- Zookeeper (matches your 3.4.10 version) --> <dependency> <groupId>org.apache.zookeeper</groupId> <artifactId>zookeeper</artifactId> <version>3.4.10</version> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> </exclusion> </exclusions> </dependency> </dependencies>
2. Core Source Code
Create a class extending AbstractPollableSource (Flume’s recommended base for polling-based sources):
import org.apache.flume.*; import org.apache.flume.conf.Configurable; import org.apache.flume.source.AbstractPollableSource; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class CustomKafkaSource extends AbstractPollableSource implements Configurable { private KafkaConsumer<String, String> kafkaConsumer; private String kafkaTopic; @Override public void configure(Context context) { Properties consumerProps = new Properties(); // Pull configs from Flume agent properties consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, context.getString("bootstrap.servers")); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, context.getString("group.id")); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // Manual commit for reliability kafkaConsumer = new KafkaConsumer<>(consumerProps); kafkaTopic = context.getString("topic"); kafkaConsumer.subscribe(Collections.singletonList(kafkaTopic)); } @Override protected Status doProcess() throws EventDeliveryException { Status status = Status.READY; Channel channel = getChannel(); Transaction tx = channel.getTransaction(); try { tx.begin(); // Poll Kafka for records ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofMillis(100)); if (records.isEmpty()) { status = Status.BACKOFF; tx.commit(); return status; } // Convert Kafka records to Flume Events for (var record : records) { Event event = EventBuilder.withBody(record.value().getBytes()); // Add useful headers (optional) event.getHeaders().put("kafka_offset", String.valueOf(record.offset())); event.getHeaders().put("kafka_partition", String.valueOf(record.partition())); channel.put(event); } tx.commit(); kafkaConsumer.commitSync(); // Commit Kafka offset only after successful channel write } catch (Exception e) { tx.rollback(); status = Status.BACKOFF; throw new EventDeliveryException("Failed to process Kafka records", e); } finally { tx.close(); } return status; } @Override protected void doStop() throws FlumeException { if (kafkaConsumer != null) { kafkaConsumer.close(); } super.doStop(); } }
Next, build a sink that writes Flume events to HDFS, with basic file rolling logic.
1. Dependencies (Add to Maven)
Add the Hadoop client dependency matching your Cloudera Hadoop version (adjust the version as needed):
<dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>2.6.0-cdh5.15.1</version> <!-- Example Cloudera Hadoop version --> <scope>provided</scope> <!-- Cloudera provides this --> </dependency>
2. Core Sink Code
Create a class extending AbstractSink:
import org.apache.flume.*; import org.apache.flume.conf.Configurable; import org.apache.flume.sink.AbstractSink; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import java.io.IOException; public class CustomHDFSSink extends AbstractSink implements Configurable { private FileSystem hdfsFs; private Path baseOutputPath; private Path currentFilePath; private FSDataOutputStream outputStream; private long fileSizeThreshold = 10 * 1024 * 1024; // 10MB roll threshold @Override public void configure(Context context) { try { Configuration hadoopConf = new Configuration(); // Load Cloudera's Hadoop configs (adjust paths as needed) hadoopConf.addResource(new Path("/etc/hadoop/conf/core-site.xml")); hadoopConf.addResource(new Path("/etc/hadoop/conf/hdfs-site.xml")); hdfsFs = FileSystem.get(hadoopConf); baseOutputPath = new Path(context.getString("hdfs.path")); // Create base directory if it doesn't exist if (!hdfsFs.exists(baseOutputPath)) { hdfsFs.mkdirs(baseOutputPath); } openNewOutputFile(); } catch (IOException e) { throw new FlumeException("Failed to configure HDFS Sink", e); } } private void openNewOutputFile() throws IOException { // Generate unique file name with timestamp String fileName = "flume-events-" + System.currentTimeMillis() + ".txt"; currentFilePath = new Path(baseOutputPath, fileName); outputStream = hdfsFs.create(currentFilePath); } @Override public Status process() throws EventDeliveryException { Status status = Status.READY; Channel channel = getChannel(); Transaction tx = channel.getTransaction(); try { tx.begin(); Event event = channel.take(); if (event == null) { status = Status.BACKOFF; tx.commit(); return status; } // Write event body to HDFS (add newline for readability) outputStream.write(event.getBody()); outputStream.write('\n'); // Check if we need to roll the file if (outputStream.getPos() >= fileSizeThreshold) { outputStream.close(); openNewOutputFile(); } tx.commit(); } catch (Exception e) { tx.rollback(); status = Status.BACKOFF; throw new EventDeliveryException("Failed to write event to HDFS", e); } finally { tx.close(); } return status; } @Override public void stop() { try { if (outputStream != null) { outputStream.close(); } if (hdfsFs != null) { hdfsFs.close(); } } catch (IOException e) { getLogger().error("Error closing HDFS resources", e); } super.stop(); } }
Critical issues you’re likely facing:
- Classpath Conflicts: Cloudera ships with its own versions of Hadoop, Kafka, and Zookeeper jars. Mark these dependencies as
<scope>provided</scope>in yourpom.xmlto avoid clashes. - Permission Issues: Ensure the Flume user (
flumeby default) has write access to your target HDFS path. Runhdfs dfs -chown flume:flume /your/hdfs/pathif needed. - Kafka Offset Compatibility: Flume 1.6’s native source uses Zookeeper for offset storage, but your custom source uses Kafka’s built-in offset management (recommended for Kafka 1.0.0). Ensure your Kafka brokers are configured to allow offset storage.
- Kerberos Authentication: If your Cloudera cluster uses Kerberos, add the JAAS config to your Flume agent startup command:
export FLUME_JAVA_OPTS="-Djava.security.auth.login.config=/path/to/flume_jaas.conf"
- Package Your Code: Use Maven to build a fat jar (or a regular jar with dependencies excluded if you marked them as provided).
- Copy Jar to Flume Lib: Place the jar in
/usr/lib/flume-ng/lib/(Cloudera’s default Flume lib directory). - Create Flume Agent Config: Write a config file pointing to your custom source/sink:
agent.sources = customKafkaSource agent.sinks = customHdfSink agent.channels = memChannel # Custom Kafka Source Config agent.sources.customKafkaSource.type = com.yourcompany.flume.CustomKafkaSource agent.sources.customKafkaSource.bootstrap.servers = your-kafka-broker:9092 agent.sources.customKafkaSource.group.id = flume-custom-group agent.sources.customKafkaSource.topic = your-input-topic # Memory Channel Config agent.channels.memChannel.type = memory agent.channels.memChannel.capacity = 10000 agent.channels.memChannel.transactionCapacity = 1000 # Custom HDFS Sink Config agent.sinks.customHdfSink.type = com.yourcompany.flume.CustomHDFSSink agent.sinks.customHdfSink.hdfs.path = hdfs://your-nn:8020/user/flume/kafka-output agent.sinks.customHdfSink.channel = memChannel
- Start the Flume Agent:
flume-ng agent -n agent -c /etc/flume-ng/conf -f /path/to/your/custom-agent.conf
内容的提问来源于stack exchange,提问作者Deepan Ram

