You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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.

实现Spark从Kafka Topic流式读取数据(Hortonworks环境)

一、补全/验证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.11 dependency—it's more intuitive for complex data pipelines

内容的提问来源于stack exchange,提问作者user9733732

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 10:38:10