spack streaming
·
spack stream简介
简单来说,Spark Streaming 是 Apache Spark 框架中用于处理“实时数据流”的组件。你可以把它理解成一个“流水线处理器”,能让数据像流水一样,一边到达、一边就被处理,而不是像传统批处理那样,攒成一大堆再统一处理。
它的核心价值是将实时数据流(无界数据)转换为一系列小的“微批次”(有界数据)进行处理。
-
传统批处理:好比快递站攒够一车货再发车,时效性差,但效率高。
-
Spark Streaming:好比快递站每隔1分钟就发一趟车,每辆车处理这1分钟内到的货。这既保证了近实时的低延迟(通常是秒级),又复用了Spark强大的批处理引擎。
常见的应用场景包括:实时网站流量统计、金融交易风控、物联网设备数据监控、实时日志分析等。
离线计算和流式计算介绍
| 计算模式 | 核心特点 | 代表框架与技术 |
|---|---|---|
| 离线计算 (Batch/Offline) | 处理静态、有界、海量的数据集,注重高吞吐量和准确性,对响应时间要求不高(通常为分钟级、小时级甚至天级)。 | Hadoop MapReduce、Hive、Spark Core |
| 流式计算 (Streaming/Real-time) | 处理动态、无界、持续产生的数据流,注重低延迟(秒级、毫秒级)和实时响应,能立即对数据进行分析和处理。 | Apache Flink、Apache Storm、Spark 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();
}
更多推荐


所有评论(0)