spack stream简介

简单来说,Spark Streaming 是 Apache Spark 框架中用于处理“实时数据流”的组件。你可以把它理解成一个“流水线处理器”,能让数据像流水一样,一边到达、一边就被处理,而不是像传统批处理那样,攒成一大堆再统一处理。

它的核心价值是将实时数据流(无界数据)转换为一系列小的“微批次”(有界数据)进行处理。

  • 传统批处理:好比快递站攒够一车货再发车,时效性差,但效率高。

  • Spark Streaming:好比快递站每隔1分钟就发一趟车,每辆车处理这1分钟内到的货。这既保证了近实时的低延迟(通常是秒级),又复用了Spark强大的批处理引擎。

常见的应用场景包括:实时网站流量统计、金融交易风控、物联网设备数据监控、实时日志分析等。

离线计算和流式计算介绍

计算模式 核心特点 代表框架与技术
离线计算 (Batch/Offline) 处理静态、有界、海量的数据集,注重高吞吐量准确性,对响应时间要求不高(通常为分钟级、小时级甚至天级)。 Hadoop MapReduceHiveSpark Core
流式计算 (Streaming/Real-time) 处理动态、无界、持续产生的数据流,注重低延迟(秒级、毫秒级)和实时响应,能立即对数据进行分析和处理。 Apache FlinkApache StormSpark Streaming

Spark Streaming-Java版本代码

导入依赖

<dependencies>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-streaming_2.12</artifactId>
        <version>3.3.1</version>
    </dependency>

    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-core_2.12</artifactId>
        <version>3.3.1</version>
</dependency>

<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-streaming-kafka-0-10_2.12</artifactId>
    <version>3.3.1</version>
 </dependency>
</dependencies>

SparkStreaming 基础环境模板

/**
 * SparkStreaming Java 版本基础环境演示代码
 * Spark对流式计算的核心运行环境、调度逻辑做了完整封装
 */
public class SparkStreaming01_Env {
    public static void main(String[] args) throws InterruptedException {
        // 1. 创建Spark核心配置对象,用于设置应用运行参数
        SparkConf conf = new SparkConf();
        // 设置运行模式为本地多线程模式,*代表自动匹配本机CPU核心数,仅本地调试使用
        conf.setMaster("local[*]");
        // 设置当前Spark应用名称,会展示在SparkUI、YARN任务列表中
        conf.setAppName("sparkStreaming");

        /**
         * 2. 创建SparkStreaming上下文对象(流式程序核心入口)
         * 参数1:Spark配置conf
         * 参数2:Duration(1000),批次间隔1000毫秒=1秒,每1秒生成一批数据执行计算
         * 该对象封装了流式任务调度、数据采集、批处理、资源管理全部核心能力
         */
        JavaStreamingContext jsc = new JavaStreamingContext(conf, new Duration(1000));

        // ====================== 缺失核心逻辑区域 ======================
        // TODO 这里必须编写:1.数据源读取  2.数据转换计算  3.输出行动算子(print/foreachRDD等)
        // 如果缺少数据源和输出算子,调用start()时Spark会校验失败,抛出 No output operations 异常
        // ===========================================================

        // 3. 启动流式任务调度器,开始持续采集、处理数据流
        // 该方法仅启动后台调度线程,不会阻塞主线程,执行完会立刻向下运行
        jsc.start();

        // 4. 阻塞主线程,让流式程序长期持续运行,不会执行完main方法就退出进程
        // 只有手动关闭程序、程序异常崩溃、主动调用stop()才会解除阻塞
        jsc.awaitTermination();
    }
}

网络数据流socket处理演示

下面我们展示了监听本机的9999端口,获取数据

public static void main(String[] args) throws InterruptedException {
        SparkConf conf = new SparkConf();
        conf.setMaster("local[*]");
        conf.setAppName("sparkStreaming");
        JavaStreamingContext jsc = new JavaStreamingContext(conf, new Duration(5000));
        //todo: 通过环境对象获取Socket数据源,获取数据模型,进行数据处理
        JavaReceiverInputDStream<String> socketDS = jsc.socketTextStream("localhost", 9999);
        //将获取的数据直接进行打印
        socketDS.print();
        jsc.start();

        //阻塞代码,他会一直阻塞在这里
        jsc.awaitTermination();
    }

使用工具

演示网络数据流的传输。使用命令进行启动,指定端口号:nc -lp 9999

kafka数据源处理

public static void main(String[] args) throws InterruptedException {
        SparkConf conf = new SparkConf();
        conf.setMaster("local[*]");
        conf.setAppName("sparkStreaming");
        JavaStreamingContext jsc = new JavaStreamingContext(conf, new Duration(5000));
        //todo: 通过kafka获取数据源,获取数据模型,进行数据处理
        HashMap<String, Object> map = new HashMap<>();
        //配置kafka的消费组配置
        map.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.1.189:9092");
        map.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        map.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        map.put(ConsumerConfig.GROUP_ID_CONFIG,"heihei");
        map.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"latest");

        // 需要消费的主题
        ArrayList<String> strings = new ArrayList<>();
        strings.add("sparkstreamingtopic");

        JavaInputDStream<ConsumerRecord<String, String>> directStream = KafkaUtils
                .createDirectStream(jsc, LocationStrategies.PreferBrokers(), ConsumerStrategies.<String, String>Subscribe(strings,map));

        directStream.map(new Function<ConsumerRecord<String, String>, String>() {
            @Override
            public String call(ConsumerRecord<String, String> v1) throws Exception {
                return v1.value();
            }
        }).print(100);

        jsc.start();

        //阻塞代码,他会一直阻塞在这里
        jsc.awaitTermination();
    }

Logo

智能硬件社区聚焦AI智能硬件技术生态,汇聚嵌入式AI、物联网硬件开发者,打造交流分享平台,同步全国赛事资讯、开展 OPC 核心人才招募,助力技术落地与开发者成长。

更多推荐