96SEO 2026-02-23 14:26 26
(xmu.edu.cn)https://dblab.xmu.edu.cn/blog/3160/2.Structured

(apache.org)https://spark.apache.org/docs/3.2.0/structured-streaming-kafka-integration.html3.Pyspark手册DataStreamReader
pyspark.sql.streaming.DataStreamReader
(apache.org)https://spark.apache.org/docs/3.2.0/api/python/reference/api/pyspark.sql.streaming.DataStreamReader.html#pyspark.sql.streaming.DataStreamReader4.Kafka安装教程
Ubuntu22.04下安装kafka_2.12-2.6.0并运行简单实例_ubuntu22.04kafka安装-CSDN博客文章浏览阅读723次点赞22次收藏31次。
安装Kafka
测试Kafka是否正常工作_ubuntu22.04kafka安装https://blog.csdn.net/qq_67822268/article/details/1386264125.Maven中央仓库
(maven.org)https://repo1.maven.org/maven2/6.Maven
mvnrepository.comhttps://mvnrepository.com/7.kafka-python文档
documentationhttps://kafka-python.read***docs.io/en/master/index.html8.strcuted
(jianshu.com)https://www.jianshu.com/p/ed1398c2470a
提供端到端的容错保证指在分布式系统中整个数据处理流程从数据输入到输出的全过程都能够保证容错性。
换句话说无论是数据的接收、处理还是输出系统都能够在发生故障或异常情况时保持数据的完整性和一致性能够确保在发生故障时不会丢失数据并且能够保证精确一次处理语义。
引擎的优化能力能够进行查询优化、状态管理和分布式处理从而提供高性能的实时处理能力。
支持事件时间event-time处理可以轻松处理乱序事件、延迟事件等场景并提供丰富的窗口操作支持。
提供了丰富的监控和调试功能包括进度报告、状态查询等方便用户监控作业的执行情况。
Streaming的关键思想是将实时数据流视为一张正在不断添加数据的表
可以把流计算等同于在一个静态表上的批处理查询Spark会在不断添加数据的无界输入表上运行计算并进行增量查询
在无界表上对输入的查询将生成结果表系统每隔一定的周期会触发对无界表的计算并更新结果表
Streaming默认使用微批处理执行模型这意味着Spark流计算引擎会定期检查流数据源并对自上一批次结束后到达的新数据执行批量查询
中偏移量Offset是指用于标识数据流中位置的标记它表示了数据流中的一个特定位置或者偏移量。
在流处理中偏移量通常用于记录已经处理的数据位置以便在失败恢复、断点续传或者状态管理等场景下能够准确地从中断处继续处理数据。
结构化流启动时会从数据源中读取偏移量并使用这些偏移量来确定应该从哪里开始读取数据。
随着数据被处理Spark
会不断更新偏移量以确保在发生故障或重启情况下能够准确地恢复到之前处理的位置。
容错和故障恢复记录偏移量可以确保在流处理过程中发生故障或者需要重启时能够准确地恢复到之前处理的位置避免数据的丢失和重复处理。
通过记录偏移量流处理系统能够知道从哪里继续读取数据从而保证数据处理的完整性和一致性。
精确一次处理语义记录偏移量也有助于实现精确一次处理语义即确保每条输入数据只被处理一次。
通过准确记录偏移量并在发生故障后能够准确地恢复到之前的位置流处理系统能够避免重复处理数据从而确保处理结果的准确性。
断点续传记录偏移量还使得流处理系统能够支持断点续传的功能即在流处理过程中可以随时停止并在之后恢复到之前的处理位置而不需要重新处理之前已经处理过的数据。
通过记录偏移量结构化流处理可以实现精确一次处理语义并确保即使在出现故障和重启的情况下也能够保证数据不会被重复处理或丢失。
因此偏移量在结构化流处理中扮演着非常重要的角色是实现流处理的容错性和准确性的关键之一。
驱动程序通过将当前待处理数据的偏移量保存到预写日志中来对数据处理进度设置检查点以便今后可以使用它来重新启动或恢复查询。
Re-executions和端到端语义在下一个微批处理之前就要将该微批处理所要处理的数据的偏移范围保存到日志中。
所以当前到达的数据需要等待先前的微批作业处理完成且它的偏移量范围被记入日志后才能在下一个微批作业中得到处理这会导致数据到达和得到处理并输出结果之间的延时超过100毫秒。
微批处理的数据延迟对于大多数实际的流式工作负载如ETL和监控已经足够了然而一些场景确实需要更低的延迟。
比如在金融行业的信用卡欺诈交易识别中需要在犯罪分子盗刷信用卡后立刻识别并阻止但是又不想让合法交易的用户感觉到延迟从而影响用户的使用体验这就需要在
1020毫秒的时间内对每笔交易进行欺诈识别这时就不能使用微批处理模型而需要使用持续处理模型。
Spark从2.3.0版本开始引入了持续处理的试验性功能可以实现流计算的毫秒级延迟。
在持续处理模式下Spark不再根据触发器来周期性启动任务而是启动一系列的连续读取、处理和写入结果的长时间运行的任务。
为了缩短延迟引入了新的算法对查询设置检查点在每个任务的输入数据流中一个特殊标记的记录被注入。
当任务遇到标记时任务把处理后的最后偏移量异步任务的执行不必等待其他任务完成或某个事件发生地报告给引擎引擎接收到所有写入接收器的任务的偏移量后写入预写日志。
由于检查点的写入是完全异步的任务可以持续处理因此延迟可以缩短到毫秒级。
也正是由于写入是异步的会导致数据流在故障后可能被处理超过一次以上所以持续处理只能做到“至少一次”的一致性。
因此需要注意到虽然持续处理模型能比微批处理模型获得更好的实时响应性能但是这是以牺牲一致性为代价的。
微批处理可以保证端到端的完全一致性而持续处理只能做到“至少一次”的一致性。
processing将连续的数据流按照一定的时间间隔或者数据量划分成小批量进行处理每个批量数据被视为一个微批作业类似于批处理的方式进行处理。
持续处理continuous
processing对不间断的数据流进行实时处理没有明确的批次边界数据到达后立即进行处理和输出。
微批处理通常会导致一定的延迟因为数据需要等待下一个批次的处理才能输出结果因此微批处理一般无法做到完全的实时性。
持续处理具有更好的实时性因为数据到达后立即进行处理可以更快地输出结果。
微批处理通常通过检查点机制来实现容错和状态管理每个微批作业之间会保存处理状态以便故障恢复和重新执行。
持续处理也需要考虑容错和状态管理但通常需要使用更复杂的机制来实现实时的状态管理和故障恢复。
微批处理可以更好地利用批处理系统的资源因为可以对数据进行分批处理适用于一些需要大批量数据一起处理的场景。
持续处理需要更多的实时资源和更高的实时性能适用于对数据要求实时性较高的场景。
Streaming采用的数据抽象是DStream本质上就是一系列RDD而Structured
Streaming采用的数据抽象是DataFrame。
Structured
SQL的DataFrame/Dataset来处理数据流。
虽然Spark
Streaming可以处理结构化的数据流。
这样Structured
Streaming可以对DataFrame/Dataset应用各种操作包括select、where、groupBy、map、filter、flatMap等。
Spark
Streaming只能实现秒级的实时响应而Structured
Streaming由于采用了全新的设计方式采用微批处理模型时可以实现100毫秒级别的实时响应采用持续处理模型时可以支持毫秒级的实时响应。
导入pyspark模块创建SparkSession对象创建输入数据源定义流计算过程启动流计算并输出结果
实例任务一个包含很多行英文语句的数据流源源不断到达Structured
Streaming程序对每行英文语句进行拆分并统计每个单词出现的频率
在/home/hadoop/sparksj/mycode/structured目录下创建StructuredNetworkWordCount.py文件
StructuredNetworkWordCountspark
\.appName(StructuredNetworkWordCount)
WARN减少日志输出spark.sparkContext.setLogLevel(WARN)#
lines.select(explode(split(lines.value,
等待流查询终止query.awaitTermination()在执行StructuredNetworkWordCount.py之前需要启动HDFS
/home/hadoop/sparksj/mycode/structured
执行程序后在“数据源终端”内用键盘不断敲入一行行英文语句nc程序会把这些数据发送给StructuredNetworkWordCount.py程序进行处理
输出结果内的Batch后面的数字说明这是第几个微批处理系统每隔8秒会启动一次微批处理并输出数据。
如果要停止程序的运行则可以在终端内键入“CtrlC”来停止。
File源或称为“文件源”以文件流的形式读取某个目录中的文件支持的文件格式为csv、json、orc、parquet、text等。
需要注意的是文件放置到给定内打开文件写入内容而是应当采取大部分操作系统都支持的、通过写入到临时文件后移动文件到给定目录的方式来完成。
path输入路径的或glob通配符路径的格式不支持以多个逗号分隔的形式。
maxFilesPerTrigger每个触发器中要处理的最大新文件数默认无最大值。
latestFirst是否优先处理最新的文件当有大量文件积压时设置为True可以优先处理新文件默认为False。
fileNameOnly是否仅根据文件名而不是完整路径来检查新文件默认为False。
如果设置为True则以下文件将被视为相同的文件因为它们的文件名“dataset.txt”相同
特定的文件格式也有一些其他特定的选项具体可以参阅Spark手册内DataStreamReader中的相关说明
以一个JSON格式文件的处理来演示File源的使用方法主要包括以下两个步骤
创建程序生成JSON格式的File源测试数据创建程序对数据进行统计
文件。
模拟了用户的登录、登出和购买行为包括事件发生的时间戳、动作类型和地区等信息
在/home/hadoop/sparksj/mycode/structured目录下创建a.py文件
检查最终及其内容os.mkdir(TEST_DATA_DIR)
程序的入口如果作为脚本直接执行则会执行下面的代码test_setUp()
清理测试环境删除最终接着使用for循环一千次来生成一千个文件文件名为“e-mall-数字.json”
其中时间、操作和省与地区均随机生成。
测试数据是模拟电子商城记录用户的行为可能是登录、退出或者购买并记录了用户所在的省与地区。
为了让程序运行一段时间每生成一个文件后休眠1秒。
在临时。
同样在/home/hadoop/sparksj/mycode/structured目录下创建b.py文件
中导入结构类型和时间戳类型、字符串类型TEST_DATA_DIR_SPARK
StructType([StructField(eventTime,
定义事件时间字段类型为时间戳StructField(action,
定义行为字段类型为字符串StructField(district,
SparkSession如果已存在则获取否则创建一个新的spark
\.appName(StructuredEMallPurchaseCount)
设置应用程序名称.getOrCreate()spark.sparkContext.setLogLevel(WARN)
每次触发处理的最大文件数以控制处理速度.load(TEST_DATA_DIR_SPARK)windowDuration
对购买行为进行筛选、按地区和时间窗口进行分组统计购买次数并按时间窗口排序windowedCounts
控制台输出不截断.trigger(processingTime10
触发处理的时间间隔.start()query.awaitTermination()
等待查询终止该程序的目的是过滤用户在电子商城里的购买记录并根据省与地区以1分钟的时间窗口统计各个省与地区的购买量并按时间排序后输出。
/home/hadoop/sparksj/mycode/structured
/home/hadoop/sparksj/mycode/structured
意思就是处理时间触发器的批处理已经开始滞后。
具体来说当前批处理花费的时间超过了触发器设定的时间间隔
毫秒也就是10秒但是当前批处理花费了16341毫秒远远超过了设定的时间间隔
当批处理花费的时间超过触发器设定的时间间隔时可能会导致处理延迟因为下一个批处理可能无法按时启动。
如果批处理持续花费较长时间可能会导致资源如CPU、内存等的浪费因为资源被用于等待而不是实际的处理任务。
上述警告可通过修改b.py代码中processingTime的值将它改成大于上图中的16341ms即可1秒1000毫秒
将b.py代码中的spark.sparkContext.setLogLevel(WARN)改为spark.sparkContext.setLogLevel(ERROR)即可
文件。
它模拟了用户的登录、登出和购买行为包括事件发生的时间戳、动作类型和地区等信息。
应用程序用于实时处理模拟的电商购买行为数据。
它从指定的读取数据并进行实时统计计算每个地区在一分钟内的购买次数并按时间窗口排序然后将结果输出到控制台。
联系a.py生成的模拟购买行为数据是b.py的输入数据源。
a.py生成的
如果你先执行a.py生成了购买行为的模拟数据然后再执行b.py它将会从a.py生成的目录中读取数据并进行实时统计购买行为数据。
这样你就可以通过实时监控控制台输出了解每个地区在一分钟内的购买情况从而进行实时的业务分析或监控。
assign指定所消费的Kafka主题和分区。
subscribe订阅的Kafka主题为逗号分隔的主题列表。
subscribePattern订阅的Kafka主题正则表达式可匹配多个主题。
kafka.bootstrap.serversKafka服务器的列表逗号分隔的“hostport”列表。
startingOffsets起始位置偏移量。
endingOffsets结束位置偏移量。
failOnDataLoss布尔值表示是否在Kafka
数据可能丢失时主题被删除或位置偏移量超出范围等触发流计算失败。
一般应当禁止以免误报。
实例使用生产者程序每0.1秒生成一个包含2个字母的单词并写入Kafka的名称为“wordcount-topic”的主题Topic内。
Spark的消费者程序通过订阅wordcount-topic会源源不断收到单词并且每隔8秒钟对收到的单词进行一次词频统计把统计结果输出到Kafka的主题wordcount-result-topic内同时通过2个监控程序检查Spark处理的输入和输出结果。
新建一个终端记作“Zookeeper终端”输入下面命令启动Zookeeper服务不要关闭这个终端窗口一旦关闭Zookeeper服务就停止了
./bin/zookeeper-server-start.sh
另外打开第二个终端记作“Kafka终端”然后输入下面命令启动Kafka服务不要关闭这个终端窗口一旦关闭Kafka服务就停止了
再新开一个终端记作“监控输入终端”执行如下命令监控Kafka收到的文本
./bin/kafka-console-consumer.sh
再新开一个终端记作“监控输出终端”执行如下命令监控输出的结果文本
./bin/kafka-console-consumer.sh
在/home/hadoop/sparksj/mycode/structured/kafkasource目录下创建并编辑spark_ss_kafka_producer.py文件
/home/hadoop/sparksj/mycode/structured/kafkasource
KafkaProducer(bootstrap_servers[localhost:9092])while
(random.choice(string.ascii_lowercase)
秒producer.send(wordcount-topic,
在运行生产者程序之前要先安装kafka-python如果读者之前已经安装可跳过此小节。
1.首先确认有没有安装pip3如果没有使用如下命令安装笔者已经安装不在演示
/home/hadoop/sparksj/mycode/structured/kafkasource
生产者程序执行以后在“监控输入终端”的窗口内就可以看到持续输出包含2个字母的单词。
程序会生成随机字符串并将其发送到
./bin/kafka-console-consumer.sh
的控制台消费者同时运行生产者程序时生产者代码会不断地生成随机字符串并发送到
主题而控制台消费者则会从该主题中读取并显示这些消息。
因此会导致生产者不断地生成消息并且控制台消费者会即时地输出这些消息从而实现了消息的生产和消费过程。
环境的搭建和消息传递的过程以确保生产者能够成功地将消息发送到指定的主题同时消费者能够从该主题中接收并处理这些消息。
同样在/home/hadoop/sparksj/mycode/structured/kafkasource目录下创建并编辑spark_ss_kafka_consumer.py文件
/home/hadoop/sparksj/mycode/structured/kafkasource
\.appName(StructuredKafkaWordCount)
设置日志级别为WARN避免过多的输出信息spark.sparkContext.setLogLevel(WARN)#
指定数据源格式为Kafka.option(kafka.bootstrap.servers,
从Kafka主题中加载数据.selectExpr(CAST(value
创建一个流式DataFrame.outputMode(complete)
指定输出数据源格式为Kafka.option(kafka.bootstrap.servers,
指定输出的Kafka主题.option(checkpointLocation,
设置检查点目录.trigger(processingTime8
等待流式查询终止在运行消费者程序即spark_ss_kafka_consumer.py时请确保kafka成功启动监控输入终端与监控输出端成功启动生产者程序成功启动若采用方式一启动消费者程序则可以等会生产者程序因为jar包下载可能时间过长长时间生产者程序会产生大量的数据若采用方式二启动消费者程序则确保启动消费者程序前启动生产者程序正如下方视频所示
org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0
使用了--packages参数指定了要从Maven仓库中下载并包含的依赖包其中org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0是要添加的Kafka相关依赖。
作用在运行应用程序时动态下载Kafka相关的依赖包并将其添加到类路径中以便应用程序能够访问这些依赖
运行后会解析包依赖并从Maven中心仓库下载所需的JAR包下载完成后进行运行但这种方法依赖于自身网络环境笔者这边因为是校园网贼慢故不再展示运行结果
在执行下列代码之前需要下载spark-sql-kafka-0-10_2.12-3.2.0.jar、kafka-clients-2.6.0.jar、commons-pool2-2.9.0.jar和spark-token-provider-kafka-0-10_2.12-3.2.0.jar文件笔者spark版本为spark
3.2.0、kafka版本为kafka_2.12-2.6.0读者请根据自己的版本调整jar版本的下载将其放到“/usr/local/spark/jars”目录下现附上下载地址
spark-sql-kafka-0-10_2.12-3.2.0.jar文件下载页面https://mvnrepository.com/artifact/org.apache.spark/spark-sql-kafka-0-10_2.12/3.2.0
kafka-clients-2.6.0.jar文件下载页面https://mvnrepository.com/artifact/org.apache.kafka/kafka-clients/2.6.0
commons-pool2-2.9.0.jar文件下载页面https://mvnrepository.com/artifact/org.apache.commons/commons-pool2/2.9.0
spark-token-provider-kafka-0-10_2.12-3.2.0.jar文件下载页面https://mvnrepository.com/artifact/org.apache.spark/spark-token-provider-kafka-0-10_2.12/3.2.0
若上述网站不能打开可尝试电脑连接手机热点或使用如下网址进行下载
链接https://pan.baidu.com/s/121zVsgc4muSt9rgCWnJZmw
spark-sql-kafka-0-10_2.12-3.2.0.jar文件下载页面
org/apache/spark/spark-sql-kafka-0-10_2.12/3.2.0
kafka-clients-2.6.0.jar文件下载页面Central
org/apache/kafka/kafka-clients/2.6.0
commons-pool2-2.9.0.jar文件下载页面Central
org/apache/commons/commons-pool2/2.9.0
spark-token-provider-kafka-0-10_2.12-3.2.0.jar文件下载页面Central
org/apache/spark/spark-token-provider-kafka-0-10_2.12/3.2.0
/usr/local/kafka/libs/*:/usr/local/spark/jars/*
使用了--jars参数指定了要包含在类路径中的外部JAR包的路径
/usr/local/kafka/libs/*和/usr/local/spark/jars/*是要包含的Kafka和Spark相关的JAR包的路径
运行如下所示同样可以设置输出日志级别来控制日志的输出在此不再赘述
./bin/kafka-console-consumer.sh
wordcount-topic在终端监控名为wordcount-topic的Kafka主题的输入信息
./bin/kafka-console-consumer.sh
在终端监控名为wordcount-result-topic的Kafka主题的输出信息
生成随机的两个小写字母字符串并将其发送到wordcount-topic主题中
从wordcount-topic主题中读取消息对单词进行计数然后将结果写入wordcount-result-topic主题。
该程序会持续运行并等待新的输入消息
监控输入终端会显示从wordcount-topic主题中接收到的随机小写字母字符串监控输入终端会显示从wordcount-result-topic主题中接收到的单词计数结果生产者程序会不断地生成随机字符串并将其发送到wordcount-topic主题消费者程序会持续地从wordcount-topic主题中读取消息对单词进行计数并将结果写入wordcount-result-topic主题
如果只执行第一条命令和生产者程序那么会看到终端不断打印出随机的两个小写字母字符串而不会有单词计数或结果输出。
IP地址或者域名必须设置。
port端口号必须设置。
includeTimestamp是否在数据行内包含时间戳。
使用时间戳可以用来测试基于时间聚合的功能。
Socket源从一个本地或远程主机的某个端口服务上读取数据数据的编码为UTF8。
因为Socket源使用内存保存读取到的所有数据并且远端服务不能保证数据在出错后可以使用检查点或者指定当前已处理的偏移量来重放数据所以它无法提供端到端的容错保障。
Socket源一般仅用于测试或学习用途。
Rate源可每秒生成特定个数的数据行每个数据行包括时间戳和值字段。
时间戳是消息发送的时间值是从开始到当前消息发送的总个数从0开始。
Rate源一般用来作为调试或性能基准测试。
rowsPerSecond每秒产生多少行数据默认为1。
rampUpTime生成速度达到rowsPerSecond
需要多少启动时间使用比秒更精细的粒度将会被截断为整数秒默认为0秒。
numPartitions使用的分区数默认为Spark的默认分区数。
源会尽可能地使每秒生成的数据量达到rowsPerSecond可以通过调整numPartitions以尽快达到所需的速度。
这几个参数的作用类似一辆汽车从0加速到100千米/小时并以100千米/小时进行巡航的过程通过增加“马力”numPartitions可以使得加速时间rampUpTime更短。
在/home/hadoop/sparksj/mycode/structured/ratesource目录下新建文件spark_ss_rate.py
\.appName(TestRateStreamSource)
设置日志级别为WARNspark.sparkContext.setLogLevel(WARN)#
等待流处理的终止query.awaitTermination()
/home/hadoop/sparksj/mycode/structured/ratesource
输出的第一行即上图红框框住的那一行StruckType就是print(lines.schema)输出的数据行的格式。
当运行这段代码时它会生成模拟的连续数据流并将其写入控制台进行显示。
输出结果会包含时间戳和生成的值。
同时程序会持续运行直到手动终止或出现异常。
同4处理警告也可以设置日志输出等级来忽略警告将spark.sparkContext.setLogLevel(WARN)改为spark.sparkContext.setLogLevel(ERROR)
作为专业的SEO优化服务提供商,我们致力于通过科学、系统的搜索引擎优化策略,帮助企业在百度、Google等搜索引擎中获得更高的排名和流量。我们的服务涵盖网站结构优化、内容优化、技术SEO和链接建设等多个维度。
| 服务项目 | 基础套餐 | 标准套餐 | 高级定制 |
|---|---|---|---|
| 关键词优化数量 | 10-20个核心词 | 30-50个核心词+长尾词 | 80-150个全方位覆盖 |
| 内容优化 | 基础页面优化 | 全站内容优化+每月5篇原创 | 个性化内容策略+每月15篇原创 |
| 技术SEO | 基本技术检查 | 全面技术优化+移动适配 | 深度技术重构+性能优化 |
| 外链建设 | 每月5-10条 | 每月20-30条高质量外链 | 每月50+条多渠道外链 |
| 数据报告 | 月度基础报告 | 双周详细报告+分析 | 每周深度报告+策略调整 |
| 效果保障 | 3-6个月见效 | 2-4个月见效 | 1-3个月快速见效 |
我们的SEO优化服务遵循科学严谨的流程,确保每一步都基于数据分析和行业最佳实践:
全面检测网站技术问题、内容质量、竞争对手情况,制定个性化优化方案。
基于用户搜索意图和商业目标,制定全面的关键词矩阵和布局策略。
解决网站技术问题,优化网站结构,提升页面速度和移动端体验。
创作高质量原创内容,优化现有页面,建立内容更新机制。
获取高质量外部链接,建立品牌在线影响力,提升网站权威度。
持续监控排名、流量和转化数据,根据效果调整优化策略。
基于我们服务的客户数据统计,平均优化效果如下:
我们坚信,真正的SEO优化不仅仅是追求排名,而是通过提供优质内容、优化用户体验、建立网站权威,最终实现可持续的业务增长。我们的目标是与客户建立长期合作关系,共同成长。
Demand feedback