首页 文章 精选 留言 我的

精选列表

搜索[流计算],共10006篇文章
优秀的个人博客,低调大师

Apache Flink 1.12.1 发布,流处理框架

Apache Flink 1.12 系列的首个 bug 修复版本 1.12.1 已经发布。该版本包含 79 个修复和优化,因此官方强烈建议所有用户都升级到 1.12.1。 Maven 依赖 <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>1.12.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.11</artifactId> <version>1.12.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients_2.11</artifactId> <version>1.12.1</version> </dependency> 注意事项 Apache Flink 1.12.1 的 DockerHub 官方映像暂时丢失。作为替代,这些映像目前放在 Flink PMC 的托管仓库中。这也是目前 Native Kubernetes 部署的默认设置。 Flink PMC 将继续与 DockerHub 团队合作以提供官方映像。 由于项目空间限制,PyPI 上暂时缺少 Apache Flink 1.12.1 的源代码和 python 3.8 linux wheel 软件包。目前,有关增加空间限制的请求正在 PyPI 审核过程中。在这段时间内,用户可以根据需要手动构建软件包。 部分更新内容 Sub-task 添加有关 maxwell-json 格式的文档 重做命令行接口文档页面 重做PyFlink CLI 文档 Bug BlobClientTest.testGetFailsDuringStreamingForJobPermanentBlob 挂起 使用 Class.forName 时加载不同的驱动程序类时出现死锁 由于超时而无法初始化 logger:引发 LoggerInitializationException 修复 ignore-parse-errors 不适用于旧版 JSON 格式 ZooKeeper quorum 因缺少 log4j 库而无法开始 Improvement 日志开始/结束状态恢复 将 “Flink Architecture” 页面翻译成中文 默认情况下启用 log4j2 监视间隔 详情请查看更新公告。

优秀的个人博客,低调大师

Apache Carbondata接入Kafka实时流数据

1.导入carbondata依赖的jar包 将apache-carbondata-1.5.3-bin-spark2.3.2-hadoop2.7.2.jar导入$SPARKHOME/jars;或将apache-carbondata-1.5.3-bin-spark2.3.2-hadoop2.7.2.jar导入在$SPARKHOME创建的carbondlib目录 2.导入kafka依赖的jar包 接入kafka数据需要依赖kafka的jars,将以下jars导入$SPARKHOME/jars kafka-clients-0.10.0.1.jarspark-sql-kafka-0-10_2.11-2.3.2.jar 3.spark-shell启动服务 ./bin/spark-shell --master spark://hostname:7077 --jars apache-carbondata-1.5.3-bin-spark2.3.2-hadoop2.7.2.jar a).导入依赖 import org.apache.spark.sql.SparkSession import org.apache.spark.sql.CarbonSession._ b).创建session 启动第一个目录是数据存储目录,第二个目录是元数据目录;都可以是hdfs目录 val carbon = SparkSession.builder().config(sc.getConf).getOrCreateCarbonSession("/home/bigdata/carbondata/data","/home/bigdata/carbondata/carbon.metastore") c).创建source表 carbon.sql( s""" | CREATE TABLE IF NOT EXISTS kafka_json_source( | id STRING, | name STRING, | age INT, | brithday TIMESTAMP) | STORED AS carbondata | TBLPROPERTIES( | 'streaming'='source', | 'format'='kafka', | 'kafka.bootstrap.servers'='hostname:9092', | 'subscribe'='kafka_json', | 'record_format'='json', | 'comment'='get kafka data') """.stripMargin).show() d).创建sink表 carbon.sql( s""" | CREATE TABLE IF NOT EXISTS kafka_json_sink( | id STRING, | name STRING, | age INT, | brithday TIMESTAMP) | STORED AS carbondata | TBLPROPERTIES( | 'streaming'='sink') """.stripMargin).show() e).创建job任务 carbon.sql( s""" | CREATE STREAM kafka_json_job ON TABLE kafka_json_sink( | STMPROPERTIES( | 'trigger'='ProcessingTime', | 'interval'='10 seconds') | AS SELECT * FROM kafka_json_source """.stripMargin).show() f).创建DATAMAP carbon.sql( s""" | CREATE DATAMAP agg_kafka_json_sink | ON TABLE kafka_json_sink( | USING "preaggregate" | AS | SELECT id,name,sum(age),max(age),min(age),avg(age) | FROM kafka_json_sink | GRPUP BY id,name """.stripMargin).show() 4.常用SQL命令 a).导入本地数据 carbon.sql("LOAD DATA INPATH '/home/bigdata/carbondata/sample.csv' INTO TABLE kafka_json_source").show() b).查看表结构 carbon.sql("DESC kafka_json_source").show() c).查看表数据 carbon.sql("SELECT * FROM kafka_json_source WHERE id=1").show() d).清理表数据 carbon.sql("TRUNCATE TABLE kafka_json_sink").show() e).删除表 carbon.sql("DROP TABLE IF EXISTS kafka_json_source").show() f).查看job任务状态 carbon.sql("SHOW STREAMS ON TABLE kafka_json_sink").show() g).删除job任务 carbon.sql("DROP STREAM kafka_json_job").show() h).查询DATAMAP表信息 carbon.sql("DESC agg_kafka_json_sink_kafka_json_sink").show() i).查询表Segments信息 carbon.sql("SHOW SEGMENTS FOR TABLE kafka_json_sink").show() j).条件查询 carbon.sql("SELECT * FROM kafka_json_sink WHERE agent_id=499 AND signature=''").show() k).聚合查询 carbon.sql("SELECT agent_id,signature,method_type,sum(elapse_time),max(elapse_time),min(elapse_time) FROM kafka_json_sink GROUP BY agent_id,signature,method_type").show() 5.注意事项 a).kafka使用配置 由于Carbondata的kafka-consumer反序列化配置如下,所以在kafka-producer应该使用对于配置,否则无法解析数据 key.deserializer = org.apache.kafka.common.serialization.ByteArrayDeserializer value.deserializer = org.apache.kafka.common.serialization.ByteArrayDeserializer

资源下载

更多资源
Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

用户登录
用户注册