在Hortonworks中用Spark Streaming从Kafka取数遇连接超时异常求助
Hey there! Since you've already got Zookeeper and Kafka up and running in your Hortonworks environment, plus a Maven project set up in Eclipse with all the streaming dependencies sorted, let's walk through exactly how to stream data from a Kafka topic using Spark.
一、补全/验证Maven POM依赖
You mentioned you've built your Maven project with streaming dependencies, here's a complete, HDP-compatible POM configuration to ensure all required components are included (note version alignment with your Hortonworks Data Platform version):
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.yourorg.spark</groupId> <artifactId>spark-kafka-streaming-demo</artifactId> <version>1.0-SNAPSHOT</version> <dependencies> <!-- Spark Core (match HDP's Spark version) --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.11</artifactId> <version>2.4.8</version> <!-- HDP 3.1 uses Spark 2.4.x --> <scope>provided</scope> </dependency> <!-- Spark Streaming --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming_2.11</artifactId> <version>2.4.8</version> <scope>provided</scope> </dependency> <!-- Spark-Kafka Integration (0.10 API for better compatibility) --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.11</artifactId> <version>2.4.8</version> </dependency> </dependencies> <build> <plugins> <!-- Compile with Java 8 (HDP standard) --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.8.1</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> <!-- Build executable shaded JAR --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.2.4</version> <executions> <execution> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>com.yourorg.spark.KafkaStreamProcessor</mainClass> </transformer> </transformers> </configuration> </execution> </executions> </plugin> </plugins> </build> </project>
Pro tip: Double-check that your Spark/Kafka versions match what's bundled in your HDP cluster—mismatched versions are the #1 cause of weird runtime errors.
二、Spark Streaming代码示例(Scala)
Here's a straightforward implementation to read from your Kafka topic, process (in this case, print) the data, and manage offsets properly:
package com.yourorg.spark import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} import org.apache.spark.streaming.kafka010.{ConsumerStrategies, KafkaUtils, LocationStrategies} object KafkaStreamProcessor { def main(args: Array[String]): Unit = { // Initialize Spark config val conf = new SparkConf() .setAppName("HDP-Spark-Kafka-Reader") .setMaster("local[*]") // Remove this line for cluster deployment // Create StreamingContext with 5-second batch interval val ssc = new StreamingContext(conf, Seconds(5)) // Kafka connection parameters val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "kafka-broker-1:9092,kafka-broker-2:9092", // Replace with your brokers "key.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "value.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "group.id" -> "spark-streaming-consumer-group", "auto.offset.reset" -> "latest", // Start from newest if no offset exists "enable.auto.commit" -> (false: java.lang.Boolean) // Manual offset management ) // Target Kafka topic(s) val topics = Array("your-target-topic") // Create direct stream to read Kafka data val kafkaStream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, // Distribute partitions evenly across workers ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ) // Process incoming data (example: print key/value pairs) kafkaStream.foreachRDD { rdd => rdd.foreach { record => println(s"Received: Key = ${record.key()}, Value = ${record.value()}") } // Manually commit offsets to ensure exactly-once processing kafkaStream.asInstanceOf[org.apache.spark.streaming.kafka010.CanCommitOffsets] .commitAsync(rdd.asInstanceOf[org.apache.spark.streaming.kafka010.HasOffsetRanges].offsetRanges) } // Start the streaming job and wait for termination ssc.start() ssc.awaitTermination() } }
三、Run the Job on Hortonworks Cluster
Once you've built your shaded JAR, submit it to the HDP cluster using spark-submit:
spark-submit --class com.yourorg.spark.KafkaStreamProcessor --master yarn --deploy-mode cluster your-shaded-jar-file.jar
Make sure:
- Your Spark cluster can reach the Kafka brokers (check firewall rules for port 9092)
- The consumer group ID is unique to avoid conflicts with other consumers
- For structured data processing, consider using Spark Structured Streaming with the
spark-sql-kafka-0-10_2.11dependency—it's more intuitive for complex data pipelines
内容的提问来源于stack exchange,提问作者user9733732

