首页 文章 精选 留言 我的

精选列表

搜索[实时etl],共10006篇文章
优秀的个人博客,低调大师

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

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

如何使用GraphDB做商品实时推荐

几乎所有的企业都需要了解如何快速并且高效地影响客户来购买他们的产品并且推荐其他相关商品给他们。这可能需要用到云服务的推荐,个性化,网络分析工具。图非常适合这些类似的分析用例,如推荐产品,或基于用户数据,过去行为,推荐个性化广告。下面我们来看看怎么使用图做个性化推荐 购买graphdb服务 登录www.aliyun.com 后进入hbase产品控制台,选择创建hbase集群 选择Graphdb(图)子产品,之后分别勾选付费方式,地域,可用区,网络类型,vpc。跟选购其他产品一致 接来下选择master及core的规则,购买core个数及磁盘容量,参考自身实际需求,这里笔者这里选择4核8G,之后选择立即购买 之后同意协议完成支付,回到hbase控制台,当前状态初始化,等待集群创建完毕 集群创建完毕后,可以拿到如下图库地址,记住这个地址,替换下面命

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

face_recognition 实时人脸识别

目标 识别进入摄像头的人是谁 face_recognition face_recognition 是github上一个非常有名气的人脸识别开源工具包,我们可以通过以下指令安装到python环境内 $ pip install face_recognition 代码的设计思路 加载认识的人脸图 ray_image = face_recognition.load_image_file("ray.jpg") ray_face_encoding = face_recognition.face_encodings(ray_image)[0] 将认识的人脸变量加到数组内 known_face_encodings = [ ray_face_encoding ] known_face_names = [ "Ray" ] 全部代码如下所示: import face_recognition import cv2 video_capture = cv2.VideoCapture(0) ray_image = face_recognition.load_image_file("ray.jpg") ray_face_encoding = face_recognition.face_encodings(ray_image)[0] pinky_image = face_recognition.load_image_file("pinky.jpg") pinky_face_encoding = face_recognition.face_encodings(ray_image)[0] # Create arrays of known face encodings and their names known_face_encodings = [ ray_face_encoding, pinky_face_encoding ] known_face_names = [ "Ray", "Pinky" ] # Initialize some variables face_locations = [] face_encodings = [] face_names = [] process_this_frame = True while True: # Grab a single frame of video ret, frame = video_capture.read() # Resize frame of video to 1/4 size for faster face recognition processing small_frame = cv2.resize(frame, (0, 0), fx=0.25, fy=0.25) # Convert the image from BGR color (which OpenCV uses) to RGB color (which face_recognition uses) rgb_small_frame = small_frame[:, :, ::-1] # Only process every other frame of video to save time if process_this_frame: # Find all the faces and face encodings in the current frame of video face_locations = face_recognition.face_locations(rgb_small_frame) face_encodings = face_recognition.face_encodings(rgb_small_frame, face_locations) face_names = [] for face_encoding in face_encodings: # See if the face is a match for the known face(s) matches = face_recognition.compare_faces(known_face_encodings, face_encoding) name = "Unknown" # If a match was found in known_face_encodings, just use the first one. if True in matches: first_match_index = matches.index(True) name = known_face_names[first_match_index] face_names.append(name) process_this_frame = not process_this_frame # Display the results for (top, right, bottom, left), name in zip(face_locations, face_names): # Scale back up face locations since the frame we detected in was scaled to 1/4 size top *= 4 right *= 4 bottom *= 4 left *= 4 # Draw a box around the face cv2.rectangle(frame, (left, top), (right, bottom), (0, 0, 255), 2) # Draw a label with a name below the face cv2.rectangle(frame, (left, bottom - 35), (right, bottom), (0, 0, 255), cv2.FILLED) font = cv2.FONT_HERSHEY_DUPLEX cv2.putText(frame, name, (left + 6, bottom - 6), font, 1.0, (255, 255, 255), 1) # Display the resulting image cv2.imshow('Video', frame) # Hit 'q' on the keyboard to quit! if cv2.waitKey(1) & 0xFF == ord('q'): break # Release handle to the webcam video_capture.release() cv2.destroyAllWindows()

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

Nacos

Nacos

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

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等操作系统。

用户登录
用户注册