首页 文章 精选 留言 我的

精选列表

搜索[玩转Lighthouse],共1675篇文章
优秀的个人博客,低调大师

基于AliOS Things玩转智能语音

随着AI技术的进步,智能语音开始将人机交互从手+眼睛的传统模式中解放出来。带给人们更便捷、更风趣、更有人情味的体验,让被操作对象变得不再只是一个死板的工具,而更像是一个有生命的助理。“帮我打开空调”,“明天上班需要带伞吗”,“帮我冲100块钱话费”…在万物互联的时代,你的所有需求只需要一句话便能实现。AliOS Things 集成的Link Voice SDK即可实现智能语音交互。 关于阿里智能语音服务 阿里智能语音服务为设备提供语音交互能力、丰富的音乐内容、智能家居控制等,并可进行专有设备技能定制(如:语音操控跑步机、按摩椅等设备)。包括: 通用服务:搜歌、搜栏目、搜电台、问天气、百科、四则运算等; 阿里服务:控制智能家居、充值手机费、天猫超市购物、查询电费等 (需接入账号体系,可参考SDS接入); 私有服务:操控设备、售后电话查询等 (

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

Android开发实践:玩转EditText控件

Android的EditText控件是一个非常常用的控件,用得最多的地方就是做登录、注册页面了,它能为用户提供一个直观便捷的输入框。本文简单总结下EditText控件中比较常用的一些设置,并为每一种设置提供两种方式的实现,一种是在布局文件中实现,另一种是在程序中通过代码动态的设置。 1. 如何添加一个方框 在Android的Hololight主题下,EditText控件默认是只有一条底部的蓝色横线的,怎么给你的EditText添加一个方框呢? 【布局】: 设置 android:background 属性,给它一个长方形的白***片,或者自定义一个长方形的drawable文件即可。 例如: 1 android:background= "@drawable/shape_bg" 【代码】: 1 2 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setBackgroundResource(R.drawable.shape_bg); 2. 如何设置字体大小、颜色、加粗 【布局】: 布局中的属性依次为 android:textSize,android:textColor,android:textStyle属性 例如: 1 2 3 android:padding="15sp" android:textSize="15sp" android:textStyle="bold" 【代码】: 1 2 3 4 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setTextSize( 15 ); mEditText.setTextColor(Color.BLACK); mEditText.setTypeface(Typeface.DEFAULT_BOLD); 3. 如何设置以密码的形式显示 【布局】: 设置 android:password 属性为 true 例如: 1 android:password="true" 【代码】: 1 2 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setInputType(InputType.TYPE_TEXT_VARIATION_PASSWORD); 4. 如何禁止用户输入回车换行 【布局】: 设置 android:singleLine 属性为 true 例如: 1 android:singleLine="true" 【代码】: 1 2 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setSingleLine(); 5. 如何设置没有输入时的提示信息 【布局】: 设置 android:hint 属性的值 例如: 1 android:hint="inputyourname" 【代码】: 1 2 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setHint( "Inputyourname" ); 6. 如何在输入框的行首空几个字符 【布局】: 设置 android:paddingLeft 属性即可 例如: 1 android:paddingLeft="15sp" 【代码】: 1 2 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setPadding( 15 , 0 , 0 , 0 ); 7. 如何限制输入的长度 【布局】: 设置 android:maxLength 属性的值即可 例如: 1 android:maxLength="10" 【代码】: 1 2 3 4 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); InputFilter[]filters= new InputFilter[ 1 ]; filters[ 0 ]= new InputFilter.LengthFilter( 10 ); mEditText.setFilters(filters); 8. 如何限制输入类型为:数字,电话号码,日期,时间 【布局】: 设置 android:inputType 属性可以指定 textPassword, phone, number, date,time 等类型 例如: 1 android:inputType="text" 【代码】: 1 2 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setInputType(InputType.TYPE_CLASS_TEXT); //InputType有很多种类型可以选择 9. 如何限制只能输入指定的字符 【布局】: 设置 android:digits 属性即可 例如: 1 android:digits="abcdef" 【代码】: 有两种方法可以实现: 方法一: 1 2 3 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); Stringdigits= "abcdef" ; mEditText.setKeyListener(DigitsKeyListener.getInstance(digits)); 方法二: 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); InputFilter[]filters= new InputFilter[ 1 ]; filters[ 0 ]= new MyInputFilter( "abcdef" ); mEditText.setFilters(filters); public class MyInputFilter extends LoginFilter.UsernameFilterGeneric{ private StringmAllowedDigits; public PopInputFilter(Stringdigits){ mAllowedDigits=digits; } @Override public boolean isAllowed( char c){ if (mAllowedDigits.indexOf(c)!=- 1 ){ return true ; } return false ; } } 10. 让密码的输入字体大小与明文的字体一致 当你设置了android:password = "true" 属性后,你会发现,它的字体大小会跟没有设置password属性的EditText的大小不一致,因此,如果期望他们表现一致的话,可以通过代码如下设置: 1 2 3 EditTextmEditText=(EditText)findViewById(R.id.MyEditText); mEditText.setTypeface(Typeface.DEFAULT); mEditText.setTransformationMethod( new PasswordTransformationMethod()); 本文转自 Jhuster 51CTO博客,原文链接:http://blog.51cto.com/ticktick/1333414,如需转载请自行联系原作者

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

玩转Android monkey之多参数实战

monkey设置间隔时间 参数throttle用来控制执行速度,如果不加我们通过上次的执行发现速度比较快,也看不清。 语法:adb shell monkey -p 包名 --throttle 设置每次随机事件的时间间隔 (毫秒)随机事件次数 moneky seed种子 seed种子是干嘛的呢?很好理解,如果你想使得两次命令的执行轨迹一样,那就加上这个参数。比如,第一次你使用发现了一个bug,现在想重现一下,如果不加seed就是随机的,可能无法重现,加上seed就可以了。 PS:每次执行时初始界面要一致 语法:adb shell monkey -p 包名 --throttle 执行速度 -s seed种子 随机事件数 monkey指定某个动作 如果想使用monkey命令来做某一个动作,而不是N个动作混杂的,那就要通过参数来指定ta的动作,比如--pct-touch点击(触摸)动作。 语法:adb shell monkey -v -p 包名 --pct-touch 执行touch动作的百分比 随机事件次数 其中-v代表查看详细的结果,我们可以看到0代表touch百分比为100%执行,其余为0%。 思考:如果不加--pct-touch执行该命令会是什么样的结果呢? 这里大家可能会问到1-10代表啥呢?其实代表的是不同的操作动作,这里来list一下: 1:手势 --pct-motion 2:缩放 --pct-pinchzoom 3:轨迹球 --pct-trackball 4:屏幕旋转 --pct-rotation 5:基本导航事件,比如手机上的上、下、左、右的操作 --pct-nav 6:主导航事件,比如返回键、菜单键 --pct-majornav 7:系统导航事件,比如手机上的home键、拨号键、音量键等 --pct-syskeys 8:切换activity --pct-appswitch 9: 键盘翻转事件,举个场景就知道了,类似点击输入框,键盘弹起,点击其他区域,键盘收起 --pct-flip 10:其他事件 --pct-anyevent monkey忽略崩溃和超时 为什么要有着两个参数呢?很简单,我们在使用app的时候经常会出现超时、卡死的状况,一旦出现这样的情况,monkey是不知道怎么办的!所以,需要我们给他指令才行, 一般就是给两个参数,忽略超时和忽略崩溃。 l 忽略超时参数:--ignore-timeouts l 忽略崩溃(异常)参数:--ignore-crashes 语法:adb shell monkey -v -p 包名 --pct-touch 100 --ignore-timeouts --ignore-crashes 随机事件次数 PS:在实际操作过程中除了上述两种情况外,可能还会出现ANR的问题,如果出现那就要找到对应的log,然后交给开发去解决 本文转自 小强测试帮 51CTO博客,原文链接:http://blog.51cto.com/xqtesting/2050469,如需转载请自行联系原作者

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

特性快闪:使用 Databend 玩转 Iceberg

作者:尚卓燃(PsiACE) 澳门科技大学在读硕士,Databend 研发工程师实习生 Apache OpenDAL(Incubating) Committer https://github.com/PsiACE 几周前,Databricks 和 Snowflake 召开了各自的年度大会,除了今年一路持续走红的 AI ,数据湖/数据仓库技术的发展仍然值得关注,毕竟数据才是基本盘。Apache Iceberg 无疑是数据湖方案的大赢家,Databricks 新推出的 UniForm 为以 Apache Iceberg 和 Hudi 表格式读取 Delta 中的数据提供了进一步的支持。而 Snowflake 也适时推出了 Iceberg Tables 更新,宣称要进一步打破数据孤岛。 Databend 最近几个月正在推动的重要新特性之一就是支持读取 Apache Iceberg 表格式的数据,尽管还没有完全落地,但已经取得了不错的进展。 今天这篇文章旨在为大家提前演示这一新特性 —— 使用 Databend 挂载并查询 Iceberg Catalog ,我们将介绍 Iceberg、表格式的一些核心概念,并且试图介绍 Databend 的解决方案(包括 Databend 的多源数据目录能力和以 Rust 从头实现的 IceLake)。当然,作为演示,我们将会提供完整的 workshop ,供大家尝鲜体验。 Apache Iceberg 时至今日,越来越多的数据进入云端,并且存储在对象存储之中,但这并不能完全适应现代分析的需求。这里有两个问题需要解决:第一个是数据以何种形式组织,也就是说,如何得到更结构化的数据存储。第二个问题还要更进一步,如何为用户提供更广泛的一致性保证以及业务中需要的模式信息,以及更多适应现代分析负载的高级特性。 数据湖往往会关注并解决第一个问题,而表格式则会致力于为第二个问题提供解决方案。 Apache Iceberg 是一种高性能的开放表格式,专为大规模分析工作负载而设计,简单而又可靠。同时支持 Spark、Trino、Flink、Presto、Hive 和 Impala 等查询引擎,并且具备模式演变(Full Schema Evolution)、时间旅行和回滚等杀手级特性。另外,Apache Iceberg 的数据分片和明确定义的数据结构还使得对数据源进行并发访问更加安全、可靠和方便。 如果你对 Iceberg 感兴趣,我们也推荐阅读像 Docker, Spark, and Iceberg: The Fastest Way to Try Iceberg! 这样的文章进行探索。 表格式初探 表格式(Table Format)是一种利用文件集合存储数据的规范。它主要包含以下三个部分的定义: 如何将数据存储在文件中 如何存储相关文件的元数据 如何存储有关表本身的元数据 表格式的文件通常存储在 HDFS、S3 或 GCS 这样的底层存储服务中,上层则会对接 Databend、Snowflake 等数据仓库。相比 CSV 或 Parquet ,表格式提供了表形式的标准的结构化数据定义,无需加载到数据仓库中就可以使用。 尽管表格式领域还有像 Delta Lake 和 Apache Hudi 这样的强劲对手,但这篇文章是关于 Apache Iceberg 的,所以,还是让我们把目光转向 Apache Iceberg ,一起了解一下它的底层文件组织结构。 上图中的 s0 、s1 代表的是表的快照信息(snapshot),也就是表在某个时刻的状态。每次 commit 都会生成一个快照,每个快照都会对应一个清单列表(manifest list),而每个清单列表可以维护多个清单文件(manifest file)的地址与统计信息,包括路径和分区范围等。清单文件中会记录当前操作生成数据文件(data file)的地址和统计信息,比如列中的最大值最小值和数据行数等。 Databend 多源数据目录 要想在 Databend 中实现 Iceberg 集成,头一件是 Databend 的多源数据目录能力。多源数据目录将会允许将原本由其他数据分析系统所管理的数据挂载到 Databend 。 从设计之初,Databend 的目标就是成为云原生的 OLAP 数据仓库,并考虑到多源数据处理的问题。Databend 中的数据按三层进行组织:catalog -> database -> table ,catalog 作为数据最大一层,会包含所有的数据库和表。 团队在此基础上设计并实现对 Hive 和 Iceberg 数据目录的支持,提供配置文件和 CREATE CATALOG 语句多种挂载形式,从而支持对相关数据进行查询。 要想挂载数据位于 S3 中的 Iceberg Catalog,只需要执行下面的 SQL 语句: CREATE CATALOG iceberg_ctl TYPE=ICEBERG CONNECTION=( URL='s3://warehouse/path/to/db' AWS_KEY_ID='admin' AWS_SECRET_KEY='password' ENDPOINT_URL='your-endpoint-url' ); IceLake - Apache Iceberg 的纯 Rust 实现 尽管 Rust 生态中近年来涌现出不少数据库、大数据分析相关的新项目,但 Rust 生态中仍然缺乏成熟的 Apache Iceberg 绑定,这为 Databend 集成 Iceberg 制造了不少困难。 Databend Labs 支持并发起的 IceLake 旨在填补这一空白,并致力于建立一个开放生态系统: 用户可以从 任何 存储服务(如 s3、gcs、azblob、hdfs 等)读写 Iceberg 表。 任何 数据库都可以集成 icelake,以支持读写 Iceberg 表。 提供原生的 arrow 格式互转换的能力。 提供多种语言绑定,使其他语言可以享有 Rust 核心带来的 Iceberg 生态支持。 当前 IceLake 已经支持读取 Apache Iceberg 存储服务中的数据(Parquet 格式)。而 Databend 的 Iceberg 数据目录能力正是由 IceLake 支撑的,其设计与实现在和 Databend 集成中得到了验证。 此外,我们还与 Iceberg 社区成员携手发起并参与 iceberg-rust 项目,旨在将 icelake 中 iceberg 相关的实现贡献给上游,目前第一个版本正紧锣密鼓的开发中,欢迎关注https://github.com/apache/iceberg-rust 。 Workshop:体验 Databend 的 Iceberg 能力 在这个 Workshop 中,我们将会展示如何准备 Iceberg 表格式的数据,并以 Catalog 的形式将其挂载到 Databend 上,并执行一些基本的查询。相关的文件和配置可以在 PsiACE/databend-workshop 中找到。 如果你本身有一些符合 Iceberg 表格式的数据存放在 OpenDAL 支持的存储服务中,我们更推荐使用 Databend Cloud ,这样你就可以跳过繁琐的服务部署和数据准备流程,轻松上手 Iceberg Catalog 。 启动服务 为了简化 Iceberg 的服务部署和数据准备问题,我们将会使用 Docker 和 Docker Compose ,你需要先安装这些组件,然后编写 docker-compose.yml 文件。 version: "3" services: spark-iceberg: image: tabulario/spark-iceberg container_name: spark-iceberg build: spark/ networks: iceberg_net: depends_on: - rest - minio volumes: - ./warehouse:/home/iceberg/warehouse - ./notebooks:/home/iceberg/notebooks/notebooks environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 ports: - 8888:8888 - 8080:8080 - 10000:10000 - 10001:10001 rest: image: tabulario/iceberg-rest container_name: iceberg-rest networks: iceberg_net: ports: - 8181:8181 environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 - CATALOG_WAREHOUSE=s3://warehouse/ - CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO - CATALOG_S3_ENDPOINT=http://minio:9000 minio: image: minio/minio container_name: minio environment: - MINIO_ROOT_USER=admin - MINIO_ROOT_PASSWORD=password - MINIO_DOMAIN=minio networks: iceberg_net: aliases: - warehouse.minio ports: - 9001:9001 - 9000:9000 command: ["server", "/data", "--console-address", ":9001"] mc: depends_on: - minio image: minio/mc container_name: mc networks: iceberg_net: environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 entrypoint: > /bin/sh -c " until (/usr/bin/mc config host add minio http://minio:9000 admin password) do echo '...waiting...' && sleep 1; done; /usr/bin/mc rm -r --force minio/warehouse; /usr/bin/mc mb minio/warehouse; /usr/bin/mc policy set public minio/warehouse; tail -f /dev/null " networks: iceberg_net: 在上述的配置文件中,我们使用 MinIO 作为底层存储,Iceberg 提供表格式能力,至于 spark-iceberg ,可以帮助我们准备一些预置数据并执行转换操作。 接下来,我们在 docker-compose.yml 文件对应的目录下启动所有服务: docker-compose up -d 数据准备 在这个 Workshop 中,我们计划使用 NYC Taxis 数据集(纽约出租车搭乘数据),在 spark-iceberg 中已经内置了 Parquet 数据,我们只需要将其转化为 Iceberg 格式。 首先启用 pyspark-notebook : docker exec -it spark-iceberg pyspark-notebook 接下来我们就可以在 http://localhost:8888 使用 Jupyter Notebook : 这里我们需要运行一小段程序,实施数据转换的操作: df = spark.read.parquet("/home/iceberg/data/yellow_tripdata_2021-04.parquet") df.write.saveAsTable("nyc.taxis", format="iceberg") 第一行将会读取 Parquet 数据,而第二行将会将其转储为 Iceberg 格式。 为了验证数据是否成功转换,我们可以访问位于 http://localhost:9001 的 MinIO 实例,可以看到数据是按之前描述的 Iceberg 底层文件组织形式进行管理的。 部署 Databend 这里我们使用手动部署单节点 Databend 服务的形式,总体上部署过程可以参考 Databend 官方文档 ,需要注意的一些细节如下: 首先是需要为日志和 Meta 数据准备相关的目录 sudo mkdir /var/log/databend sudo mkdir /var/lib/databend sudo chown -R $USER /var/log/databend sudo chown -R $USER /var/lib/databend 其次,因为默认的 admin_api_address 已经被前面的服务占用掉,所以需要编辑 databend-query.toml 进行一些修改避免冲突: admin_api_address = "0.0.0.0:8088" 另外,我们还需要根据 Docs | Configuring Admin Users 配置管理员用户,由于只是一个 workshop ,这里选择最简单的方式,只是取消 [[query.users]] 字段以及 root 用户的注释: [[query.users]] name = "root" auth_type = "no_password" 由于我们本地部署 MinIO ,没有设置证书加密,需要使用不安全的 HTTP 协议加载数据,所以还需要更改 databend-query.toml 配置文件以允许这一行为。在生产服务中请尽可能避免开启它: ... [storage] ... allow_insecure = true ... 接下来就可以正常启动 Databend : ./scripts/start.sh 我们强烈推荐你使用 BendSQL 作为客户端,当然,我们也支持像 MySQL Client 和 HTTP API 等多种访问形式。 挂载 Iceberg Catalog 根据之前的配置文件,只需要执行下述 SQL 就可以一键挂载 Iceberg Catalog 。 CREATE CATALOG iceberg_ctl TYPE=ICEBERG CONNECTION=( URL='s3://warehouse/' AWS_KEY_ID='admin' AWS_SECRET_KEY='password' ENDPOINT_URL='http://localhost:9000' ); 为了验证是否成功,我们可以执行 SHOW CATALOGS 查看: 当然,我们也支持了 SHOW DATABASES 和 SHOW TABLES 语句,之前转换数据时的 nyc.taxis 对应在 MinIO 中是二级目录,而在 Databend 则会映射到数据库和表。 执行查询 数据已经挂载,那么就让我们执行一些简单的查询: 首先是对数据进行行数统计,可以看到一共挂载了 200 万行数据到 Databend: SELECT count(*) FROM iceberg_ctl.nyc.taxis; 让我们从其中几列试着取一些数据出来: SELECT tpep_pickup_datetime, tpep_dropoff_datetime, passenger_count FROM iceberg_ctl.nyc.taxis LIMIT 5; 下面的查询可以帮助我们探索乘客数量和旅程距离之间的相关性,这里只取其中 10 条结果: SELECT passenger_count, to_year(tpep_pickup_datetime) AS year, round(trip_distance) AS distance, count(*) FROM iceberg_ctl.nyc.taxis GROUP BY passenger_count, year, distance ORDER BY year, count(*) DESC LIMIT 10; 总结 在这篇文章中,我们介绍到 Apache Iceberg 表格式和 Databend 的相关解决方案,并且提供了一个相对完整的 workshop 供大家探索。 不难看出,尽管目前我们只为 Iceberg Catalog 提供了单机模式的目录挂载能力,但 Databend 可以胜任一些基本的查询处理任务。也欢迎大家在自己感兴趣的数据上进行尝试,并给我们提供一些反馈。

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

Apache APISIX 玩转 Tongsuo 国密插件

本文通过解读国密的相关内容与标准,呈现了当下国内技术环境中对于国密功能支持的现状。并从 API 网关 Apache APISIX 的角度,带来有关国密的探索与功能呈现。 作者:罗泽轩,Apache APISIX PMC 什么是国密 顾名思义,国密就是国产化的密码算法。在我们日常开发过程中会接触到各种各样的密码算法,如 RSA、SHA256 等等。为了达到更高的安全等级,许多大公司和国家会制定自己的密码算法。国密就是这样一组由中国国家密码管理局制定的密码算法。在国际形势越发复杂多变的今天,密码算法的国产化替代,在一些领域已经成为了一股势不可挡的潮流。 国密的官方名称为国家商用密码,简称商密,拼音缩写是 SM。这也是国密标准中 SM2/3/4/7/9 等算法名称的来源。国密算法的命名方式非常简单直接,就像“绵阳九所”、“二机部”一样,都是“分类+序号”的组合。其中 SM1 和 SM4 是对称算法,对标 AES;SM2 是非对称算法,对标 RSA、ECDSA;SM3 是摘要算法,对标 MD5;等等。由于本文并非涉及到国密的实现细节,所以不会讲得非常细,也就科普下国密算法的分类。 在基础的国密算法之上,我们可以构造一个国密的垂直生态,比如实现国密算法的硬件、提供国密支持的密码库、加入国密流程的 TLS 握手协议等等。正如安全需要纵深防御一样,基于国密的信任链也需要有全软件栈上的支持。 因此,当我们谈论国密支持时,并不仅仅单独指可以用某一种国密算法进行加解密,而是指嵌合入国密的生态,支持某种国密的应用场景。 国密的应用场景 作为国家密码管理局制定的密码算法,国密广泛应用于电子政务(包括国家政务通、警务通等重要领域)、信创及金融业的各个应用领域。 政府和金融的身份认证终端。依照现行有关规定,许多涉及政府和金融的身份认证终端(诸如 USBKey、智能 IC 卡、银行卡终端等)都需要提供对国密的支持。 国产开源操作系统。许多主打国产替代的开源操作系统,会提供基于国密的安全加固功能。比如龙蜥操作系统 (Anolis OS) 提到自己实现了全栈国密能力;OpenEuler 也在做国密相关的一些功能,比如基于国密数字证书扩展了 EFI 的数字签名。 信创产品。还有许多做信创生意的厂商,围绕国密推出符合相关标准的产品。例如支持使用国密算法做数字签名的 PDF 工具、支持国密接入标准的音视频软件等等。 基于国密 TLS 协议的生态。也是在日常开发中接触得最多的。譬如各种国产 CA 厂商、支持了国密TLS的许许多多密码库和浏览器,以及国密接入的VPN和网关等等。 APISIX 对国密的探索与支持 Apache APISIX 是一个动态、实时、高性能的 API 网关,提供负载均衡、动态上游、灰度发布、精细化路由、限流限速、服务降级、服务熔断、身份认证、可观测性等数百项功能。作为一个发迹于国内环境的 API 网关,APISIX 自然需要考虑下如何接入国密的生态圈。 由于国密更多地用在国内环境中,尚未被 OpenSSL 等国际主流项目所完全接纳。如果要使用国密的功能,就涉及到更换 APISIX 默认的 OpenSSL 库为其他 SSL 库。 为此,我们考察了以下的项目: GMSSL:北京大学开源的项目。原版基于 OpenSSL 1.0.2 修改而来。新的 GMSSL 3.0 基本上是重新开发,跟 OpenSSL 的目录结构差别很大。 gm-BoringSSL:个人开源项目,在 BoringSSL 上增加国密支持。已有两年未改动。 TaSSL:北京江南天安科技有限公司开源的项目。基于 OpenSSL 1.1.1 修改而来。 Tongsuo:蚂蚁集团开源的项目。基于 OpenSSL 3.0 修改而来,项目前身是 BabaSSL,现已 改名为铜锁/Tongsuo。 由于 GMSSL 3.0 并不基于 OpenSSL,即使能保证 API 兼容,也没办法确保能 100% 替换现有 OpenSSL 的行为,所以被首先排除。其次 gm-BoringSSL 疏于维护,也被排除。 在 TaSSL 和 Tongsuo 之中,我倾向于选择 Tongsuo[1]。因为 Tongsuo 在标准上拥有更强的话语权,比如 RFC 8998(TLS 1.3 中支持 SM 套件)就是由 Tongsuo 的开发者制定的。TaSSL 则是每出一个版本,就公布一个新的仓库。比如前一个 版本[2],感觉不太靠谱。 对于选择 Tongsuo,我个人存在一个顾虑,就是他目前基于 OpenSSL 3.0 的版本还没有发布正式的 Release 版本(第一个版本预计在 2023 年 2 月发布)。由于 Tongsuo 当前还没有一个固定的版本,因此社区决定先把国密相关的功能独立出来,以插件形式存在,有相关需求时可单独启用。 目前已在插件层面实现了服务端一侧国密双证书的支持,感兴趣的读者可以在官网查看 gm 插件介绍文档 插件介绍文档[3],自行完成 APISIX 的编译和对应插件的安装配置工作。当然,如果想即刻预览该插件的使用过程,也可以直接参考下文内容。 快速参考:APISIX 国密插件的使用 启用插件 插件要求 Apache APISIX 运行在编译了 Tongsuo 的 APISIX-Base 上。 首先需要安装 Tongsuo (此处我们选择编译出 Tongsuo 的动态链接库): # TODO: use a fixed release once they have created one. # See https://github.com/Tongsuo-Project/Tongsuo/issues/318 git clone https://github.com/api7/tongsuo --depth 1 pushd tongsuo ./config shared enable-ntls -g --prefix=/usr/local/tongsuo make -j2 sudo make install_sw 其次需要构建 APISIX-Base,让它使用 Tongsuo 作为 SSL 库: export OR_PREFIX=/usr/local/openresty export openssl_prefix=/usr/local/tongsuo export zlib_prefix=$OR_PREFIX/zlib export pcre_prefix=$OR_PREFIX/pcre export cc_opt="-DNGX_LUA_ABORT_AT_PANIC -I${zlib_prefix}/include -I${pcre_prefix}/include -I${openssl_prefix}/include" export ld_opt="-L${zlib_prefix}/lib -L${pcre_prefix}/lib -L${openssl_prefix}/lib64 -Wl,-rpath,${zlib_prefix}/lib:${pcre_prefix}/lib:${openssl_prefix}/lib64" ./build-apisix-base.sh 该插件默认是禁用状态,你需要将其添加到配置文件./conf/config.yaml 中才可以启用它: plugins: - ... - gm 由于 APISIX 的默认 cipher 中不包含国密 cipher,所以我们还需要在配置文件 ./conf/config.yaml 中设置 cipher: apisix: ... ssl: ... # 可按实际情况调整。错误的 cipher 会导致 “no shared cipher” 或 “no ciphers available” 报错。 ssl_ciphers: HIGH:!aNULL:!MD5 配置完成后,重新加载 APISIX,此时 APISIX 将会启用国密相关的逻辑。 测试插件 在测试插件之前,需要准备好国密双证书。Tongsuo 提供了生成 【SM2 双证书】 的 教程[4]。 在下面的例子中,我们将用到如下的证书: # 客户端加密证书和密钥 t/certs/client_enc.crt t/certs/client_enc.key # 客户端签名证书和密钥 t/certs/client_sign.crt t/certs/client_sign.key # CA 和中间 CA 打包在一起的文件,用于设置受信任的 CA t/certs/gm_ca.crt # 服务端加密证书和密钥 t/certs/server_enc.crt t/certs/server_enc.key # 服务端签名证书和密钥 t/certs/server_sign.crt t/certs/server_sign.key 此外,还需要准备 Tongsuo 命令行工具。 ./config enable-ntls -static make -j2 # 生成的命令行工具在 apps 目录下 mv apps/openssl .. 你也可以采用非静态编译的方式,不过就需要根据具体环境,自己解决动态链接库的路径问题了。以下示例展示了如何在指定域名中启用 gm 插件。 要在指定域名中启用 gm 插件的功能,需要先创建对应的 SSL 对象: #!/usr/bin/env python # coding: utf-8 import sys # sudo pip install requests import requests if len(sys.argv) <= 3: print("bad argument") sys.exit(1) with open(sys.argv[1]) as f: enc_cert = f.read() with open(sys.argv[2]) as f: enc_key = f.read() with open(sys.argv[3]) as f: sign_cert = f.read() with open(sys.argv[4]) as f: sign_key = f.read() api_key = "edd1c9f034335f136f87ad84b625c8f1" resp = requests.put("http://127.0.0.1:9180/apisix/admin/ssls/1", json={ "cert": enc_cert, "key": enc_key, "certs": [sign_cert], "keys": [sign_key], "gm": True, "snis": ["localhost"], }, headers={ "X-API-KEY": api_key, }) print(resp.status_code) print(resp.text) 然后将上面的脚本保存为 ./create_gm_ssl.py,运行以下命令: ./create_gm_ssl.py t/certs/server_enc.crt t/certs/server_enc.key t/certs/server_sign.crt t/certs/server_sign.key 输出结果如下: 200 {"key":"\/apisix\/ssls\/1","value":{"keys":["Yn... 完成上述准备后,可以使用如下命令测试插件是否启用成功: ./openssl s_client -connect localhost:9443 -servername localhost -cipher ECDHE-SM2-WITH-SM4-SM3 -enable_ntls -ntls -verifyCAfile t/certs/gm_ca.crt -sign_cert t/certs/client_sign.crt -sign_key t/certs/client_sign.key -enc_cert t/certs/client_enc.crt -enc_key t/certs/client_enc.key 其中,./openssl 是上文提到的 Tongsuo 命令行工具,9443 是 APISIX 默认的 HTTPS 端口。 如果一切正常,可以看到连接已经建立了起来,并输出如下信息: ... New, NTLSv1.1, Cipher is ECDHE-SM2-SM4-CBC-SM3 ... 禁用插件 如果不再使用此插件,可将 gm 插件从 ./conf/config.yaml 配置文件中移除,然后重启 APISIX 或者通过插件热加载的接口触发插件的卸载。 相关链接 如果你对该功能或者插件感兴趣,欢迎在随时在社区进行交流。 [1] https://github.com/Tongsuo-Project/Tongsuo [2] https://github.com/jntass/TASSL-1.1.1k [3] https://apisix.apache.org/zh/docs/apisix/next/plugins/gm/ [4] https://www.yuque.com/tsdoc/ts/sulazb 本文由博客一文多发平台 OpenWrite 发布!

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

smart-socket实战:玩转心跳消息

一、背景 在通信中设计的心跳消息,通常是为了检查网络链路是否正常。虽然TCP协议提供keep-alive机制,但需要在链路空闲2小时后才触发检测,这显然对业务非常不友好。当存在大量连接异常,而服务端却需要等2个小时后才感知到的时候,有限的系统资源会被逐渐耗尽,最终无法为新连接请求继续提供服务。 二、原理 要解决此类问题,业界的普遍做法是在应用层加入心跳机制。心跳消息可以是单向心跳也可以是双向心跳,所谓单向心跳表示由服务端或者客户端的其中一方主动发送心跳请求消息,而另一方返回响应消息(如下图)。双向心跳表示服务端与客户端相互发送心跳请求和响应。因为无论何种类型,实现方案都是一样的,本文以单向心跳为例给大家做讲解。 三、方案 心跳消息通常是周期性的发送,或者是在链路空闲一定时长后触发。如果经历几个周期后都未收到响应,则可以视为链路异常。此时可以继续尝试发送心跳,也可以执行告警并断开连接。 在 smart-socket 中我们提供了现成的心跳插件 HeartPlugin,可以很方便的实现心跳。本文是假定读者朋友对 smart-socket 已有了初步的了解,所以不会涉及 smart-socket 的基础使用,重点描述如何在服务中集成心跳插件。 3.1 HeartPlugin插件概述 3.1.1 心跳策略 在HeartPlugin中有三种心跳策略可供选择,通过选择不同的构造方案确定。 HeartPlugin(int heartRate, TimeUnit timeUnit) heartRate 表示心跳消息的发送频率;timeUnit 表示 heartRate 的数值单位。例如:heartRate=3,timeUnit=TimeUnit.SECONDS,表示每 3秒钟发送一次心跳。heartRate=2000,timeUnit=TimeUnit.MILLISECONDS,表示每 2秒钟发送一次心跳。该策略为周期性发送心跳消息,无论对方是否返回响应。 HeartPlugin(int heartRate, int timeout, TimeUnit unit) 该构造方法相较前一个多出一个参数:timeout(过期时间),必须大于heartRate。如果在timeout时长内发送的心跳消息都没有收到响应消息,则视为链路异常并且该链路会被关闭,释放资源。 HeartPlugin(int heartRate, int timeout, TimeUnit timeUnit, TimeoutCallback timeoutCallback) 该构造方法支持指定超时回调策略 timeoutCallback,其实上一个构造方法就是设置了超时断链策略。如果不满足业务所需,用户可按需定义。 3.1.2 心跳的识别与触发 心跳策略确定好后,下一步就是如何去发送心跳消息,以及如何识别收到的消息是否为响应消息。在 HeartPlugin 中已经定义了这两个接口,需要开发人员去实现处理逻辑: sendHeartRequest 发送心跳。HeartPlugin 在判断某个连接需要触发心跳后,会执行该方法。用户需要在该方法中实现心跳消息的编码并输出数据。 public void sendHeartRequest(AioSession session) throws IOException{ WriteBuffer writeBuffer = session.writeBuffer(); byte[] heartBytes = "heart_req".getBytes(); writeBuffer.writeInt(heartBytes.length); writeBuffer.write(heartBytes); writeBuffer.flush(); } isHeartMessage 请求消息识别。true:表示本次收到的是心跳消息(请求/响应);false:其他业务消息,交由MessageProcessor#processor处理。 public boolean isHeartMessage(AioSession session, String msg) { //心跳请求消息,返回响应 if("heart_req".equals(msg)){ try { WriteBuffer writeBuffer = session.writeBuffer(); byte[] heartBytes = "heart_rsp".getBytes(); writeBuffer.writeInt(heartBytes.length); writeBuffer.write(heartBytes); writeBuffer.flush(); }catch (Exception e){ } return true; } //是否为心跳响应消息 return "heart_rsp".equals(msg); } 3.2 代码演示 3.2.1 服务端 public class HeartServer { private static final Logger LOGGER = LoggerFactory.getLogger(HeartServer.class); public static void main(String[] args) throws IOException { //定义消息处理器 AbstractMessageProcessor<String> processor = new AbstractMessageProcessor<String>() { @Override public void process0(AioSession<String> session, String msg) { LOGGER.info("收到客户端:{}消息:{}", session.getSessionID(), msg); } @Override public void stateEvent0(AioSession<String> session, StateMachineEnum stateMachineEnum, Throwable throwable) { switch (stateMachineEnum) { case SESSION_CLOSED: LOGGER.info("客户端:{} 断开连接", session.getSessionID()); break; } } }; //注册心跳插件:每隔1秒发送一次心跳请求,5秒内未收到消息超时关闭连接 processor.addPlugin(new HeartPlugin<String>(1, 5, TimeUnit.SECONDS) { @Override public void sendHeartRequest(AioSession session) throws IOException { WriteBuffer writeBuffer = session.writeBuffer(); byte[] heartBytes = "heart_req".getBytes(); writeBuffer.writeInt(heartBytes.length); writeBuffer.write(heartBytes); writeBuffer.flush(); } @Override public boolean isHeartMessage(AioSession session, String msg) { //心跳请求消息,返回响应 if ("heart_req".equals(msg)) { try { WriteBuffer writeBuffer = session.writeBuffer(); byte[] heartBytes = "heart_rsp".getBytes(); writeBuffer.writeInt(heartBytes.length); writeBuffer.write(heartBytes); writeBuffer.flush(); } catch (Exception e) { } return true; } //是否为心跳响应消息 if ("heart_rsp".equals(msg)) { LOGGER.info("收到来自客户端:{} 的心跳响应消息", session.getSessionID()); return true; } return false; } }); //启动服务 AioQuickServer<String> server = new AioQuickServer<>(8888, new StringProtocol(), processor); server.start(); } } 3.2.2 客户端 client_1:接受服务端的心跳消息,不做任何回应 client_2:及时响应服务端的心跳消息 public class HeartClient { private static final Logger LOGGER = LoggerFactory.getLogger(HeartClient.class); public static void main(String[] args) throws IOException, ExecutionException, InterruptedException { AbstractMessageProcessor<String> client_1_processor = new AbstractMessageProcessor<String>() { @Override public void process0(AioSession<String> session, String msg) { LOGGER.info("client_1 收到服务端消息:" + msg); } @Override public void stateEvent0(AioSession<String> session, StateMachineEnum stateMachineEnum, Throwable throwable) { LOGGER.info("stateMachineEnum:{}", stateMachineEnum); } }; AioQuickClient<String> client_1 = new AioQuickClient<>("localhost", 8888, new StringProtocol(), client_1_processor); client_1.start(); AbstractMessageProcessor<String> client_2_processor = new AbstractMessageProcessor<String>() { @Override public void process0(AioSession<String> session, String msg) { LOGGER.info("client_2 收到服务端消息:" + msg); try { if ("heart_req".equals(msg)) { WriteBuffer writeBuffer = session.writeBuffer(); byte[] heartBytes = "heart_rsp".getBytes(); writeBuffer.writeInt(heartBytes.length); writeBuffer.write(heartBytes); LOGGER.info("client_2 发送心跳响应消息"); } } catch (Exception e) { e.printStackTrace(); } } @Override public void stateEvent0(AioSession<String> session, StateMachineEnum stateMachineEnum, Throwable throwable) { LOGGER.info("stateMachineEnum:{}", stateMachineEnum); } }; AioQuickClient<String> client_2 = new AioQuickClient<>("localhost", 8888, new StringProtocol(), client_2_processor); client_2.start(); } } 3.2.3观察控制台 服务端 客户端 总结 本文围绕着心跳原理作了简单的实践分享。现实场景中如果对接的设备数量高达几万,甚至十几万,本文的心跳方案是否依旧适用,欢迎一起交流讨论。 本文涉及到的示例代码可从smart-socket仓库中下载

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

玩转kubernetes中Pod挂载公网IP

起因 在很多业务场景里,是需要kubernetes里的Pod去调用某些公网上服务的(无论这些公网服务是自建的还是其它服务提供商的)。但是通常Pod都是通过主机上SNAT的方式做出口,这样容易被误认为是某种信息抓取程序或者是类似DDoS攻击,从而容易被封闭对应的主机IP。 另外有些服务只是对某种服务提供接入,需要设置白名单,但是kuberentes集群里又跑了多种服务,这些服务都是通过主机的SNAT出去,从而比较难限定对应的服务的白名单。如果将对应的服务固定在某几个worker节点上虽然也是一个办法,但是其灵活度以及容器弹性将收到一定的限制。 是否可以提供一个方式,让对应的服务的Pod根据需要来挂载自己独立的公网IP,从而避免上述问题呢?答案是可以的,这个就是我们kubernetes容器服务中特有的网络插件Terway的一个很好的特性。 提

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

使用docker镜像玩转steam挂卡

概述 之前我写过怎么在steam上挂卡,就是下面这篇文章 https://www.bboysoul.com/2017/10/10/%E4%BD%BF%E7%94%A8ArchiSteamFarm%E5%9C%A8%E6%A0%91%E8%8E%93%E6%B4%BE%E6%8C%82%E5%8D%A1/ 不过自从我学习了docker之后,我发现没有什么是不能用一个镜像解决的,如果不能,那就两个,所以,从现在开始我要学会用docker解决任何问题,比如找女朋友。 首先说一下什么是挂卡 当你在steam里玩游戏的时候,你会发现当你玩的时间累积到一定的时间的时候,steam会奖励你一些卡,然后这些卡你可以在市场上卖,虽然卖出去的卡赚来的钱不能取出来,但是你可以买其他游戏啊。挂卡就是帮助你保持游戏的在线时间,然后赚取这些卡片。 但是问题又来了,我挂卡需要游戏,但是我没有钱买这些游戏怎么办?这个就由我这个老司机告诉你,首先没有游戏没关系,没有钱也没有关系,steam上经常会有一些游戏限免,这些游戏可以让你的游戏库加一,然后这些游戏一般都会有卡片的,接着你去关注下面这个商店,很多时候它都会送游戏 https://www.humblebundle.com/store 拿着领取到的key在steam上激活就好了,最后关注一些其他的喜加一新闻就好了 一些前提条件 首先肯定要docker啊,这个就不说了,很简单,在我的博客里搜索下就好了,其次最好使用国外的服务器挂卡,因为你懂的,中国大陆ping不通steamcommunity.com这个地址的 操作 说了这么多废话之后说下操作,首先clone下面这个仓库地址 git clone https://github.com/bboysoulcn/ArchiSteamFarm.git 之后star下这个仓库,并且follow这个很帅气的人,接着build这个镜像,输入下面命令 cd ArchiSteamFarm docker build -t bboysoul/archisteamfarm:3.3.0.3 . 注意上面3.3.0.3后面有个点 如果你不想build这个镜像呢也没有关系,直接pull就可以了 docker pull bboysoul/archisteamfarm:3.3.0.3 之后新建一个screen会话并且运行起来 screen -S steam docker run --name archisteamfarm -it bboysoul/archisteamfarm:3.3.0.3 sh -c "/usr/bin/vim /asf/config/bboysoul.json && /asf/ArchiSteamFarm" 首先会让你输入账号和密码,之后会有输入一个steam的验证码 全部输入完成之后,并且像下面这个样子 ArchiSteamFarm git:(master) docker run --name archisteamfarm -it archisteamfarm:3.3.0.3 sh -c "/usr/bin/vim /asf/config/bboysoul.json && /asf/ArchiSteamFarm" 2018-08-22 23:08:32|ArchiSteamFarm-7|INFO|ASF|InitASF() ArchiSteamFarm V3.3.0.3 (linux-x64/61c03fef-7e4e-4e04-abbf-00d089ff014c | Linux 4.14.14-041414-lowlatency #201801201219 SMP PREEMPT Sat Jan 20 12:23:20 UTC 2018) 2018-08-22 23:08:33|ArchiSteamFarm-7|INFO|ASF|InitGlobalConfigAndLanguage() ASF will attempt to use your preferred culture, but translation in that language was completed only in 0.0 %. Perhaps you could help us improve ASF translation for your language? 2018-08-22 23:08:33|ArchiSteamFarm-7|INFO|ASF|InitGlobalDatabaseAndServices() It looks like it's your first launch of the program, welcome! 2018-08-22 23:08:43|ArchiSteamFarm-7|WARN|ASF|InitGlobalDatabaseAndServices() Please review our privacy policy section on the wiki if you're concerned about what ASF is in fact doing! 2018-08-22 23:08:49|ArchiSteamFarm-7|INFO|ASF|CheckAndUpdateProgram() ASF will automatically check for new versions every 1 day. 2018-08-22 23:08:49|ArchiSteamFarm-7|INFO|ASF|CheckAndUpdateProgram() Checking for new version... 2018-08-22 23:08:51|ArchiSteamFarm-7|INFO|ASF|CheckAndUpdateProgram() Local version: 3.3.0.3 | Remote version: 3.3.0.3 2018-08-22 23:08:51|ArchiSteamFarm-7|INFO|ASF|InitializeSteamConfiguration() Initializing SteamDirectory... 2018-08-22 23:08:51|ArchiSteamFarm-7|INFO|ASF|InitializeSteamConfiguration() Success! 2018-08-22 23:08:52|ArchiSteamFarm-7|INFO|bboysoul|Start() Starting... 2018-08-22 23:08:52|ArchiSteamFarm-7|INFO|bboysoul|Connect() Connecting... 2018-08-22 23:08:53|ArchiSteamFarm-7|INFO|bboysoul|OnConnected() Connected to Steam! 2018-08-22 23:08:53|ArchiSteamFarm-7|INFO|bboysoul|OnConnected() Logging in... <bboysoul> Please enter SteamGuard auth code that was sent on your e-mail: 5888K 2018-08-22 23:09:17|ArchiSteamFarm-7|INFO|bboysoul|OnDisconnected() Disconnected from Steam! 2018-08-22 23:09:17|ArchiSteamFarm-7|INFO|bboysoul|OnDisconnected() Reconnecting... 2018-08-22 23:09:17|ArchiSteamFarm-7|INFO|bboysoul|Connect() Connecting... 2018-08-22 23:09:20|ArchiSteamFarm-7|INFO|bboysoul|OnConnected() Connected to Steam! 2018-08-22 23:09:20|ArchiSteamFarm-7|INFO|bboysoul|OnConnected() Logging in... 2018-08-22 23:09:20|ArchiSteamFarm-7|INFO|bboysoul|OnLoggedOn() Successfully logged on as 76561198422915309/bboysoulcn. 2018-08-22 23:09:20|ArchiSteamFarm-7|INFO|bboysoul|Init() Logging in to ISteamUserAuth... 2018-08-22 23:09:22|ArchiSteamFarm-7|INFO|bboysoul|Init() Success! 2018-08-22 23:09:22|ArchiSteamFarm-7|INFO|bboysoul|IsAnythingToFarm() Checking first badge page... 2018-08-22 23:09:24|ArchiSteamFarm-7|INFO|bboysoul|StartFarming() We have a total of 12 games (39 cards) left to idle (~22 hours, 30 minutes remaining)... 2018-08-22 23:09:24|ArchiSteamFarm-7|INFO|bboysoul|Farm() Chosen idling algorithm: Complex 2018-08-22 23:09:24|ArchiSteamFarm-7|INFO|bboysoul|FarmSolo() Now idling: 550 (Left 4 Dead 2) 2018-08-22 23:09:25|ArchiSteamFarm-7|INFO|bboysoul|ShouldFarm() Idling status for 550 (Left 4 Dead 2): 3 cards remaining 2018-08-22 23:09:25|ArchiSteamFarm-7|INFO|bboysoul|FarmCards() Still idling: 550 (Left 4 Dead 2) 就表示成功了,并且正在挂卡中 ctrl+a+d离开这个会话。 总结一下 如果用上docker,那么你整个刮开流程只要四步 安装docker 执行docker pull bboysoul/archisteamfarm:3.3.0.3 执行screen -S steam 执行docker run --name archisteamfarm -it bboysoul/archisteamfarm:3.3.0.3 sh -c "/usr/bin/vim /asf/config/bboysoul.json && /asf/ArchiSteamFarm" 和以前要安装各种依赖影响宿主机来说好多了 欢迎关注Bboysoul的博客www.bboysoul.comHave Fun

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

玩转kafka中的消费者

上一篇介绍了如何使用Kafka的生产者,这一篇将介绍在实际生产中如何合理地使用Kafka的消费者API。 Kafka中消费者API分为新版和旧版,本章只介绍新版,旧版的就不做介绍了。 首发于我的个人博客:http://www.janti.cn/article/kafkaconsumer 准备工作 kafka版本:2.11-1.1.1 操作系统:centos7 java:jdk1.8 有了以上这些条件就OK了,具体怎么安装和启动Kafka这里就不强调了,可以看上一篇文章。 新建一个maven工程,需要的依赖如下: <dependency> <groupId>org.apache.kafkagroupId> <artifactId>kafka_2.11artifactId> <version>1.1.1version> dependency> <dependency> <groupId>org.apache.kafkagroupId> <artifactId>kafka-clientsartifactId> <version>1.1.1version> dependency> 简单的消费者例子 Kafka中是封装了KafkaConsumer类,消息接受都是通过该类来进行的。与生产者一样,在实例化消费者之前都是要进行配置的。 首先介绍这里面的配置: bootstrap.servers:配置连接代理列表,不必包含Kafka集群的所有代理地址,当连接上一个代理后,会从集群元数据信息中获取其他存活的代理信息。但为了保证能够成功连上Kafka集群,在多代理集群的情况下,建议至少配置两个代理。(由于电脑配置有限,本文实验的是单机情况) key.deserializer: 用于反序列化消息Key的类 value.deserializer:用于反序列化消息值(Value)的类 group.id:指定消费者所在的组 client.id:指定客户端所在组的ID enable.auto.commit:设置是否自动提交。在没有指定消费偏移量提交方式时,默认是每隔1S发送一次 auto.commit.interval.ms:自动提交偏移值的时间间 订阅消息的流程分为以下: 1.消费者参数配置 2.实例化消费者 3.订阅主题 4.从Kafka中拉取消息 按照这样的流程代码如下: package kafka.consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.util.Arrays; import java.util.Properties; /** * 简单的消费者 * * @author tangj */ public class KafkaSimpleDemo { static Properties properties = new Properties(); private static String topic = "MyOrder"; //poll超时时间 private static long pollTimeout = 2000; // 1.消费者配置 static { // bootstarp server 地址 properties.put("bootstrap.servers", "10.0.90.53:9092"); // group.id指定消费者,所在的组 properties.put("group.id", "order"); // 组中client 的ID名称 properties.put("client.id", "consumer"); // 在没有指定消费偏移量提交方式时,默认是每个1s提交一次偏移量,可以通过auto.commit.interval.ms参数指定提交间隔 // 自动提交要设置成true // 手动提交设置成false properties.put("enable.auto.commit", true); // 自动提交偏移值的 时间间隔 properties.put("auto.commit.interval.ms", 2000); // key序列化 properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); // value序列化 properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); } public static void main(String args[]) { consumeAutoCommit(); } /** * 消费消息自动提交偏移值 */ public static void consumeAutoCommit() { //2\. 实例化消费者 KafkaConsumer kafkaConsumer = new KafkaConsumer(properties); // 3.订阅主题 kafkaConsumer.subscribe(Arrays.asList(topic)); try { // 消费者是一个长期的过程,所以使用永久循环, while (true) { // 4.拉取消息 ConsumerRecords records = kafkaConsumer.poll(pollTimeout); for (ConsumerRecord record : records) { System.out.println("消息总数为: " + records.count()); System.out.println("收到消息: " + String.format("partition = %d, offset= %d, key=%s, value=%s%n", record.partition(), record.offset(), record.key(), record.value())); } } } catch (Exception e) { e.printStackTrace(); } finally { kafkaConsumer.close(); } } } 看了简单的消费者模式,来看看消费者中的其他知识:消费模型,分区均衡,偏移量的提交,多线程消费者。 消费模型 可以看到在配置的时候,配置了消费者组这个变量。这里主要讲讲消费模型。消费模型主要分为两种:消费者组模型,发布订阅模型 消费者组中有多个消费者,组内的消费者共享一个group.id,并且有一个单独的client.id。消费者组用来消费一个主题下的所有分区,但是每个分区只能由该组内的一个消费者消费,不会被重复消费。 发布订阅模型就是所有的消费者都可以通过订阅来获取Kafka中的消息。 分区再平衡 介绍一种情况,消费者并不是越多越好的,当消费者大于分区数时,就会有部分消费者一直空闲着。 分区再平衡,即:parition rebanlance。总的来说分区再平衡就是一个消费者原来消费的分区变成由其他消费者消费,它只发生在消费者组中。 再均衡的作用就是为了保证消费者组的高可用和伸缩性,但是再均衡期间消费者会无法读取消息,有短暂的暂停时间。 偏移量的处理 偏移量的左右就是记录已经消费的消息。之前的一个简单的demo演示了自动提交的方式,接下来介绍手动提交。 在实际生产中,消费者拉取到消息之后会进行一些业务处理,比如存到数据库,写入缓存,网络请求等,这些都会有失败的可能,所以要对偏移值进行更精细的控制。 手动提交有两种方式: 第一种是同步提交,同步提交是阻塞的,对于提交失败的处理,它会一直提交,直到提交成功。 第二种是异步提交,异步提交是非阻塞的,对于提交失败,也不会重新提交。 当然,对于手动提交的业务设计,还是要结合具体业务进行考虑和设计。手动提交之后,需要设置 enable.auto.commit为false.并且不需要设置 auto.commit.interval.ms。 下面给出一个自动提交的例子,设置没消费5次,提交一次偏移值: public static void consumeHandleCommit() { KafkaConsumer kafkaConsumer = new KafkaConsumer(properties);try {int maxcount = 5;int count = 0; kafkaConsumer.subscribe(Arrays.asList(topic));for (; ; ) {// 拉取消息 ConsumerRecords records = kafkaConsumer.poll(polltimeout);for (ConsumerRecord record : records) { System.out.println("收到消息: " + String.format("partition = %d, offset= %d, key=%s, value=%s%n", record.partition(), record.offset(), record.key(), record.value())); count++; }// 业务逻辑完成后,提交偏移量 if (count >= maxcount) { kafkaConsumer.commitAsync(new OffsetCommitCallback() {@Override public void onComplete(Map offsets, Exception exception) {if (null != exception) { exception.printStackTrace(); } else { System.out.println("偏移值提交成功"); } } }); count = 0; } } } catch (Exception e) { e.printStackTrace(); } finally { kafkaConsumer.close(); } 多线程消费者 单线程的消费者效率肯定是低于多线程的消费者的,但是消费者的多线程设计与生产者不同,KafkaConsumer是非线程安全的。 保证线程安全,每个线程,各自实例化一个KafkaConsumer对象,并且多个消费者线程只消费同一个主题,不考虑多个消费者线程消费同一个分区。 线程实体类: package kafka.consumer;import org.apache.kafka.clients.consumer.ConsumerRecord;import org.apache.kafka.clients.consumer.ConsumerRecords;import org.apache.kafka.clients.consumer.KafkaConsumer;import java.util.Arrays;import java.util.Map;import java.util.Properties;public class KafkaConsumerThread implements Runnable {//每个线程私有一个consumer实例 private KafkaConsumer<String,String> consumer;public KafkaConsumerThread(Map<String,Object> configMap,String topic) { Properties properties = new Properties(); properties.putAll(configMap);this.consumer = new KafkaConsumer<String, String>(properties);consumer.subscribe(Arrays.asList(topic)); }@Override public void run() {try {for (; ; ) {// 拉取消息 ConsumerRecords<String, String> records = consumer.poll(1000);for (ConsumerRecord<String, String> record : records) { System.out.println("收到消息: " + String.format("partition = %d, offset= %d, key=%s, value=%s%n", record.partition(), record.offset(), record.key(), record.value())); } } } catch (Exception e) { e.printStackTrace(); } finally {consumer.close(); } } } 线程启动类: package kafka.consumer;import java.util.HashMap;import java.util.Map;import java.util.concurrent.ExecutorService;import java.util.concurrent.Executors;public class KafkaConsumerExcutor {/** * kafkaConsumer是非线程安全的,处理好多线程同步的方案是 * 每个线程实例化一个kafkaConsumer对象 */ public static void main(String args[]) { String topic = "hello"; Map<String, Object> configMap = new HashMap<>(); configMap.put("bootstrap.servers", "10.0.90.53:9092");//group.id指定消费者,所在的组,保证所有线程都在一个消费者组 configMap.put("group.id", "test"); configMap.put("enable.auto.commit", true); configMap.put("auto.commit.interval.ms", 1000);// key序列化 configMap.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");// value序列化 configMap.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); ExecutorService service = Executors.newFixedThreadPool(6);// 该主题总共有6个分区,那么保证6个线程 for (int i = 0; i < 6; i++) { service.submit(new KafkaConsumerThread(configMap, topic)); } } } 总结 本文介绍了kafka中的消费者,包括多线程消费者,偏移量的处理,以及分区再平衡。 切记一点,实际的生产中消费者需要根据实际业务来进行设计。

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

玩转阿里云Serverless Kubernetes新功能

Serverless Kubernetes产品介绍 还记得阿里云Serverless Kubernetes五月份正式对外邀测的激动时刻吗?还记得阿里云Serverless Kubernetes与传统Kubernetes集群的区别吗?让我们再次一起了解一下Serverless Kubernetes。 阿里云Serverless Kubernetes让您无需管理和维护集群与服务器,即可快速创建 Kuberentes 容器应用,并且根据应用实际使用的 CPU 和内存资源量进行按需付费。使用 Serverless Kubernetes,您可以专注于设计和构建应用程序,而不是管理运行应用程序的基础设施。它基于阿里云弹性计算基础架构,并且完全兼容 Kuberentes API的解决方案,充分结合了虚拟化资源带来的安全性、弹性和 Kubernete

资源下载

更多资源
Mario

Mario

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

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应用均可从中受益。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册