-
Process类是一个抽象类,Runtime.exec()方法可以创建一个本地进程,并返回Process子类的一个实例。 Process p = Runtime.getRuntime().exec(cmd); cmd 是字符串类型 也可以是字符串类型的数组 内容就是命令行。例如下面的代码 String[] cmds = new String[]{"python", "C:\\Users\\张家豪\\Desktop\\keras\\test1.py",url,a}; Process pcs; pcs = Runtime.getRuntime().exec(cmds); 第一个参数和第二个参数加起来就是一个命令行 python 路径\\python文件 变量1 变量2(cmd执行python脚本并传参) ,三四个参数是我代码里的要传给python程序的字符串信息。 (另外在java中,RunTime.getRuntime().exec()实现了调用服务器命令脚本来执行功能需要。 用法: public Process exec(String command)-----在单独的进程中执行指定的字符串命令。 public Process exec(String [] cmdArray)---在单独的进程中执行指定命令和变量 public Process exec(String command, String [] envp)----在指定环境的独立进程中执行指定命令和变量 public Process exec(String [] cmdArray, String [] envp)----在指定环境的独立进程中执行指定的命令和变量 public Process exec(String command,String[] envp,File dir)----在有指定环境和工作目录的独立进程中执行指定的字符串命令 public Process exec(String[] cmdarray,String[] envp,File dir)----在指定环境和工作目录的独立进程中执行指定的命令和变量) 那么,如何获取Process的数据流呢,那便是要依靠getInputStream和getOutputStream方法。 如果你要往Process进程中输入数据,那么你要调用Process的getOutputStream方法! 相反,如果你要获取Process进程的输出数据,那么你要调用Process的getInputStream方法! exitValue:返回该Process对象代表的进程的出口值,值0表示正常退出,非0非正常。关于该方法,应该是返回System.exit(int)方法中的参数。 ———————————————— 原文链接:https://blog.csdn.net/zhang2362167998/article/details/122740617
-
1.消费位移确认 Kafka消费者消费位移确认有自动提交与手动提交两种策略。在创建KafkaConsumer对象时,通过参数enable.auto.commit设定,true表示自动提交(默认)。自动提交策略由消费者协调器(ConsumerCoordinator)每隔${auto.commit.interval.ms}毫秒执行一次偏移量的提交。手动提交需要由客户端自己控制偏移量的提交。 (1)自动提交。在创建一个消费者时,默认是自动提交偏移量,当然我们也可以显示设置为自动。例如,我们创建一个消费者,该消费者自动提交偏移量 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "test"); props.put("client.id", "test"); props.put("enable.auto.commit", true);// 显示设置偏移量自动提交 props.put("auto.commit.interval.ms", 1000);// 设置偏移量提交时间间隔 props.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);// 创建消费者 consumer.subscribe(Arrays.asList("test"));// 订阅主题 (2)手动提交。在有些场景我们可能对消费偏移量有更精确的管理,以保证消息不被重复消费以及消息不被丢失。假设我们对拉取到的消息需要进行写入数据库处理,或者用于其他网络访问请求等等复杂的业务处理,在这种场景下,所有的业务处理完成后才认为消息被成功消费,这种场景下,我们必须手动控制偏移量的提交。 Kafka 提供了异步提交(commitAsync)及同步提交(commitSync)两种手动提交的方式。两者的主要区别在于同步模式下提交失败时一直尝试提交,直到遇到无法重试的情况下才会结束,同时,同步方式下消费者线程在拉取消息时会被阻塞,直到偏移量提交操作成功或者在提交过程中发生错误。而异步方式下消费者线程不会被阻塞,可能在提交偏移量操作的结果还未返 回时就开始进行下一次的拉取操作,在提交失败时也不会尝试提交。 实现手动提交前需要在创建消费者时关闭自动提交,即设置enable.auto.commit=false。然后在业务处理成功后调用commitAsync()或commitSync()方法手动提交偏移量。由于同步提交会阻塞线程直到提交消费偏移量执行结果返回,而异步提交并不会等消费偏移量提交成功后再继续下一次拉取消息的操作,因此异步提交还提供了一个偏移量提交回调的方法commitAsync(OffsetCommitCallback callback)。当提交偏移量完成后会回调OffsetCommitCallback 接口的onComplete()方法,这样客户端根据回调结果执行不同的逻辑处理。 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "test"); props.put("client.id", "test"); props.put("fetch.max.bytes", 1024);// 为了便于测试,这里设置一次fetch 请求取得的数据最大值为1KB,默认是5MB props.put("enable.auto.commit", false);// 设置手动提交偏移量 props.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); // 订阅主题 consumer.subscribe(Arrays.asList("test")); try { int minCommitSize = 10;// 最少处理10 条消息后才进行提交 int icount = 0 ;// 消息计算器 while (true) { // 等待拉取消息 ConsumerRecords<String, String> records = consumer.poll(1000); for (ConsumerRecord<String, String> record : records) { // 简单打印出消息内容,模拟业务处理 System.out.printf("partition = %d, offset = %d,key= %s value = %s%n", record. partition(), record.offset(), record.key(),record.value()); icount++; } // 在业务逻辑处理成功后提交偏移量 if (icount >= minCommitSize){ consumer.commitAsync(new OffsetCommitCallback() { @Override public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception exception) { if (null == exception) { // TODO 表示偏移量成功提交 System.out.println("提交成功"); } else { // TODO 表示提交偏移量发生了异常,根据业务进行相关处理 System.out.println("发生了异常"); } } }); icount=0; // 重置计数器 } } } catch(Exception e){ // TODO 异常处理 e.printStackTrace(); } finally { consumer.close(); } 3.5以时间戳查询消息 Kafka 在0.10.1.1 版本增加了时间戳索引文件,因此我们除了直接根据偏移量索引文件查询消息之外,还可以根据时间戳来访问消息。consumer-API 提供了一个offsetsForTimes(Map<TopicPartition, Long> timestampsToSearch)方法,该方法入参为一个Map 对象,Key 为待查询的分区,Value 为待查询的时间戳,该方法会返回时间戳大于等于待查询时间的第一条消息对应的偏移量和时间戳。需要注意的是,若待查询的分区不存在,则该方法会被一直阻塞。 假设我们希望从某个时间段开始消费,那们就可以用offsetsForTimes()方法定位到离这个时间最近的第一条消息的偏移量,在查到偏移量之后调用seek(TopicPartition partition, long offset)方法将消费偏移量重置到所查询的偏移量位置,然后调用poll()方法长轮询拉取消息。例如,我们希望从主题“stock-quotation”第0 分区距离当前时间相差12 小时之前的位置开始拉取消息 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "test"); props.put("client.id", "test"); props.put("enable.auto.commit", true);// 显示设置偏移量自动提交 props.put("auto.commit.interval.ms", 1000);// 设置偏移量提交时间间隔 props.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); // 订阅主题 consumer.assign(Arrays.asList(new TopicPartition("test", 0))); try { Map<TopicPartition, Long> timestampsToSearch = new HashMap<TopicPartition,Long>(); // 构造待查询的分区 TopicPartition partition = new TopicPartition("stock-quotation", 0); // 设置查询12 小时之前消息的偏移量 timestampsToSearch.put(partition, (System.currentTimeMillis() - 12 * 3600 * 1000)); // 会返回时间大于等于查找时间的第一个偏移量 Map<TopicPartition, OffsetAndTimestamp> offsetMap = consumer.offsetsForTimes (timestampsToSearch); OffsetAndTimestamp offsetTimestamp = null; // 这里依然用for 轮询,当然由于本例是查询的一个分区,因此也可以用if 处理 for (Map.Entry<TopicPartition, OffsetAndTimestamp> entry : offsetMap.entrySet()) { // 若查询时间大于时间戳索引文件中最大记录索引时间, // 此时value 为空,即待查询时间点之后没有新消息生成 offsetTimestamp = entry.getValue(); if (null != offsetTimestamp) { // 重置消费起始偏移量 consumer.seek(partition, entry.getValue().offset()); } } while (true) { // 等待拉取消息 ConsumerRecords<String, String> records = consumer.poll(1000); for (ConsumerRecord<String, String> record : records){ // 简单打印出消息内容 System.out.printf("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(); } 3.6消费速度控制 提供 pause(Collection<TopicPartition> partitions)和resume(Collection<TopicPartition> partitions)方法,分别用来暂停某些分区在拉取操作时返回数据给客户端和恢复某些分区向客户端返回数据操作。通过这两个方法可以对消费速度加以控制,结合业务使用。 ———————————————— 原文链接:https://blog.csdn.net/qq_35349490/article/details/79790625
-
Kafka client 消息接收的三种模式 引言 kafka的消费模式总共有3种:最多一次,最少一次,正好一次。为什么会有这3种模式,是因为客户端处理消息,提交反馈(commit)这两个动作不是原子性。 1.最多一次:客户端收到消息后,在处理消息前自动提交,这样kafka就认为consumer已经消费过了,偏移量增加。 2.最少一次:客户端收到消息,处理消息,再提交反馈。这样就可能出现消息处理完了,在提交反馈前,网络中断或者程序挂了,那么kafka认为这个消息还没有被consumer消费,产生重复消息推送。 3.正好一次:保证消息处理和提交反馈在同一个事务中,即有原子性。 本文从这几个点出发,详细阐述了如何实现以上三种方式。 1.At-most-once(最多一次) 设置enable.auto.commit为ture 设置 auto.commit.interval.ms为一个较小的时间间隔. client不要调用commitSync(),kafka在特定的时间间隔内自动提交。 示例 public void mostOnce(){ Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "test-1"); props.put("enable.auto.commit", "true"); props.put("auto.commit.interval.ms", "1000"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Arrays.asList("my-topic", "bar")); while (true) { ConsumerRecords<String, String> records = consumer.poll(100); for (ConsumerRecord<String, String> record : records) { process(record); } } } 2.At-least-once(最少一次) 方法一 设置enable.auto.commit为false client调用commitSync(),增加消息偏移; public void leastOnce(){ Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "test-1"); props.put("enable.auto.commit", "false"); //取消自动提交 props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Arrays.asList("my-topic", "bar")); while (true) { ConsumerRecords<String, String> records = consumer.poll(100); for (ConsumerRecord<String, String> record : records) process(record); consumer.commitAsync(); //提交offset } } 方法二 设置enable.auto.commit为ture 设置 auto.commit.interval.ms为一个较大的时间间隔. client调用commitSync(),增加消息偏移; 示例 public void leastOnce(){ Properties props = new Properties(); props.put("bootstrap.servers", "10.242.1.219:9092"); props.put("group.id", "test-1"); props.put("enable.auto.commit", "true"); //自动提交 props.put("auto.commit.interval.ms", "99999999"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Arrays.asList("my-topic", "bar")); while (true) { ConsumerRecords<String, String> records = consumer.poll(100); for (ConsumerRecord<String, String> record : records) process(record); consumer.commitAsync(); //提交offset } } 3.Exactly-once(正好一次) 3.1 思路 如果要实现这种方式,必须自己控制消息的offset,自己记录一下当前的offset,对消息的处理和offset的移动必须保持在同一个事务中,例如在同一个事务中,把消息处理的结果存到mysql数据库同时更新此时的消息的偏移。 3.2 实现 设置enable.auto.commit为false 保存ConsumerRecord中的offset到数据库 当partition分区发生变化的时候需要rebalance,有以下几个事件会触发分区变化 1 consumer订阅的topic中的分区大小发生变化 2 topic被创建或者被删除 3 consuer所在group中有个成员挂了 4 新的consumer通过调用join加入了group 此时 consumer通过实现ConsumerRebalanceListener接口,捕捉这些事件,对偏移量进行处理。 consumer通过调用seek(TopicPartition, long)方法,移动到指定的分区的偏移位置。 3.3 实验 首先需要建立存储topic,partition中的offset记录,建表如下 DROP TABLE IF EXISTS `tb_yx_message`; CREATE TABLE `tb_yx_message` ( `id` bigInt(20) NOT NULL AUTO_INCREMENT COMMENT '主键' , `topic` varchar(128) NOT NULL DEFAULT '' COMMENT '主题' , `kPartition` varchar(128) NOT NULL DEFAULT '0' COMMENT '分区' , `offset` bigInt(20) NOT NULL DEFAULT '' COMMENT '偏移' , PRIMARY KEY (`id`), UNIQUE KEY `UNIQ_KEY`(topic,kPartition) )ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='分区消息表';#约10万 2.同时创建操作数据库的的MessageDao类 public interface MessageDao { @Insert(" insert into tb_yx_message(topic,kPartition,offset) " + " values(#{topic},#{kPartition},#{offset})") public int insertWinner(MessageOffsetPO MessageOffset); //获取偏移 @Select(" select * from tb_yx_message " + " where topic=#{topic} and kPartition=#{kPartition} ") public MessageOffsetPO get(@Param("topic") String topic, @Param("kPartition") String partition); //更新偏移 @Update(" update tb_yx_message set offset=#{offset}" + " where topic=#{topic} and kPartition=#{kPartition}") public int update(@Param("offset") long offset, @Param("topic") String topic, @Param("kPartition") String partition); } 3.接下去需要实现ConsumerRebalanceListener接口,在分区rebalance的时候,调用的顺序:先调用onPartitionsRevoked(通知consumer 任务被取消了),再调用onPartitionsAssigned(通知consumer新的任务来了)。 4.那么在我们收到任务被取消的时候,把对应offset保存到数据库;在收到新任务到来的时候,从数据库读出对应分区的偏移(例如刚启动),具体实现如下所示。 public class MyConsumerRebalancerListener implements org.apache.kafka.clients.consumer.ConsumerRebalanceListener { private Consumer<String, String> consumer; private MessageDao messageDao; public MyConsumerRebalancerListener(Consumer<String, String> consumer, MessageDao messageDao) { this.consumer = consumer; this.messageDao = messageDao; } //任务被取消 public void onPartitionsRevoked(Collection<TopicPartition> partitions) { for (TopicPartition partition : partitions) { long offset = consumer.position(partition); MessageOffsetPO build = MessageOffsetPO .builder() .offset(offset) .kPartition(partition.partition() + "") .topic(partition.topic()) .build(); try { messageDao.insertWinner(build); } catch (Exception e) { } log.info("onPartitionsRevoked topic:{},build:{}",partition.topic(),build); } } //收到新任务 public void onPartitionsAssigned(Collection<TopicPartition> partitions) { for (TopicPartition partition : partitions) { MessageOffsetPO messageOffsetPO = messageDao.get(partition.topic(), partition.partition() + ""); if(messageOffsetPO==null){ //接收到新的topic,即数据库中记录不存在 consumer.seek(partition,0); }else{ consumer.seek(partition,messageOffsetPO.getOffset()+1);//下一个offset,所以需要加1 } log.info("onPartitionsAssigned topic:{},messageOffsetPO:{},offset:{}",partition,messageOffsetPO); } } } 对于client端的代码可以这么写 /** * 正好一次 */ public void exactlyOnce(){ Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "test-1"); props.put("enable.auto.commit", "false"); //取消自动提交 props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); MyConsumerRebalancerListener rebalancerListener = new MyConsumerRebalancerListener(consumer,messageDao); consumer.subscribe(Arrays.asList("test-new-topic-1", "new-topic"),rebalancerListener); while (true) { ConsumerRecords<String, String> records = consumer.poll(100); for (ConsumerRecord<String, String> record : records) { boolean isException=false; try { log.info("consume record:{}",record); processService.process(record); } catch (Exception e) { e.printStackTrace(); isException=true; } if(isException){ //处理发生异常,说明数据没有被消费,为了保证能被消费,需要移动到该位置继续进行处理 TopicPartition topicPartition=new TopicPartition(record.topic(),record.partition()); consumer.seek(topicPartition,record.offset()); log.info("consume exception offset:{}",record.offset()); break; } // rebalancerListener.getOffsetMananger().saveOffsetInExternalStore(record.topic(),record.kPartition(),record.offset()); } } } 第19行 **processService.process(record);**对消息进行消费,同时记录了消息的偏移位置,ProcessService代码如下 @Service @Slf4j public class ProcessService { @Autowired MessageDao messageDao; @Transactional(rollbackFor = Exception.class) public void process(ConsumerRecord<String, String> record){ log.info(" record:{}",record); //对消息进行处理,这里只是简单的打印了一下 System.out.println(">>>>>>>>>"+Thread.currentThread().getName()+"_"+record); //更新偏移量 messageDao.update(record.offset(),record.topic(),record.partition()+""); } } 3.4 运行程序. 场景1: consumer订阅topic——test-new-topic-1(之前未订阅过),然后producer向consumer发送两条信息,consumer收到信息,未抛出异常。得到日志如下: 2018-01-11 14:36:07,948 main INFO MyConsumerRebalancerListener.onPartitionsAssigned:53 -onPartitionsAssigned topic:test-new-topic-1-0,messageOffsetPO:null,offset:{} 2018-01-11 14:37:15,156 main INFO KafkaConsumeTest.exactlyOnce:106 -consume record:ConsumerRecord(topic = test-new-topic-1, partition = 0, offset = 0, CreateTime = 1515652635056, serialized key size = 1, serialized value size = 14, headers = RecordHeaders(headers = [], isReadOnly = false), key = 0, value = message:0-test) 2018-01-11 14:37:15,183 main INFO ProcessService.process:23 - record:ConsumerRecord(topic = test-new-topic-1, partition = 0, offset = 0, CreateTime = 1515652635056, serialized key size = 1, serialized value size = 14, headers = RecordHeaders(headers = [], isReadOnly = false), key = 0, value = message:0-test) 2018-01-11 14:37:15,238 main INFO KafkaConsumeTest.exactlyOnce:106 -consume record:ConsumerRecord(topic = test-new-topic-1, partition = 0, offset = 1, CreateTime = 1515652635062, serialized key size = 1, serialized value size = 14, headers = RecordHeaders(headers = [], isReadOnly = false), key = 1, value = message:1-test) 2018-01-11 14:37:15,239 main INFO ProcessService.process:23 - record:ConsumerRecord(topic = test-new-topic-1, partition = 0, offset = 1, CreateTime = 1515652635062, serialized key size = 1, serialized value size = 14, headers = RecordHeaders(headers = [], isReadOnly = false), key = 1, value = message:1-test) 第1行,consumer收到新的任务指派,由于数据库中没有这个topic,所以获得的messgeOffsetPO为null,那么consumer需要从offset为0处开始接受数据。 第2-3行,consumer获得了数据并对数据进行处理,消息的offset=0 第4-5行,consumer获得了数据并对数据进行处理,消息的offset=1 此时数据库的状态为: 场景2:producer向consumer发送两条信息,consumer收到信息,进行处理,但抛出异常。 @Transactional(rollbackFor = Exception.class) public void process(ConsumerRecord<String, String> record){ log.info(" record:{}",record); System.out.println(">>>>>>>>>"+Thread.currentThread().getName()+"_"+record); messageDao.update(record.offset(),record.topic(),record.partition()+""); throw new RuntimeException("error");//抛出异常 } 运行后得到如下结果: 2018-01-11 14:56:48,357 main INFO KafkaConsumeTest.exactlyOnce:116 -consume exception offset:2 2018-01-11 14:56:48,866 main INFO KafkaConsumeTest.exactlyOnce:106 -consume record:ConsumerRecord(topic = test-new-topic-1, partition = 0, offset = 2, CreateTime = 1515653807138, serialized key size = 1, serialized value size = 14, headers = RecordHeaders(headers = [], isReadOnly = false), key = 0, value = message:0-test) 2018-01-11 14:56:48,867 main INFO ProcessService.process:23 - record:ConsumerRecord(topic = test-new-topic-1, partition = 0, offset = 2, CreateTime = 1515653807138, serialized key size = 1, serialized value size = 14, headers = RecordHeaders(headers = [], isReadOnly = false), key = 0, value = message:0-test) ....... 可以看到,consumer一直在消费offset为2的数据,由于处理时时异常状态,所以一直在消费2,数据库此时的状态与之前的一致,offset为1,符合预期。 4.引用: kafkaClinet官方文档 kafkaProducer官方文档 老外的博客 他是用文件读写的方式实现“正好一次”,感觉不是很好,不过对我启发比较大。 ———————————————— 原文链接:https://blog.csdn.net/laojiaqi/article/details/79034798
-
1.Kafka是什么 Kafka是由Apache软件基金会开发的一个开源流处理平台,由Scala和Java编写。Kafka是一种高吞吐量的分布式发布订阅消息系统,它可以处理消费者在网站中的所有动作流数据。 这种动作(网页浏览,搜索和其他用户的行动)是在现代网络上的许多社会功能的一个关键因素。 这些数据通常是由于吞吐量的要求而通过处理日志和日志聚合来解决。 对于像Hadoop一样的日志数据和离线分析系统,但又要求实时处理的限制,这是一个可行的解决方案。Kafka的目的是通过Hadoop的并行加载机制来统一线上和离线的消息处理,也是为了通过集群来提供实时的消息。简单的说,Kafka是由Linkedin开发的一个分布式的消息队列系统(Message Queue)。kafka的架构师jay kreps非常喜欢franz kafka,觉得kafka这个名字很酷,因此将linkedin的消息传递系统命名为完全不相干的kafka,没有特别含义。 2.解决什么问题 kafka开发的主要初衷目标是构建一个用来处理海量日志,用户行为和网站运营统计等的数据处理框架。在结合了数据挖掘,行为分析,运营监控等需求的情况下,需要能够满足各种实时在线和批量离线处理应用场合对低延迟和批量吞吐性能的要求。从需求的根本上来说,高吞吐率是第一要求,其次是实时性和持久性。 既有的消息队列框架或者对消息传送的可靠性提供了较高的保证,由此带来较大的负担,不能满足海量高吞吐率的要求;或者完全面向实时消息处理系统,对于批量离线处理的场合无法提供足够的缓存和持久性要求。 而多数针对大数据开发应用的日志收集处理系统则通常更适合批量离线处理场合,对实时在线处理的场合支持不够。 总体而言,kafka试图提供一个同时满足在线和离线处理海量数据的消息派发系统。 3.怎么实现 kafka的集群由多个Broker服务器组成,每个类型的消息被定义为topic,同一topic内部的消息按照一定的key和算法被分区(partition)存储在不同的Broker上,消息生产者producer和消费者consumer可以在多个Broker上生产/消费topic ———————————————— 原文链接:https://blog.csdn.net/weixin_33567029/article/details/112076328
-
Springboot2.0整合Kafka,从Kafka并发、批量获取数据 Kafka安装 Spring Boot是由Pivotal团队提供的全新框架,其设计目的是用来简化新Spring应用的初始搭建以及开发过程。该框架使用了特定的方式来进行配置,从而使开发人员不再需要定义样板化的配置。通过这种方式,Spring Boot致力于在蓬勃发展的快速应用开发领域(rapid application development)成为领导者。Spring框架是Java平台上的一种开源应用框架,提供具有控制反转特性的容器。尽管Spring框架自身对编程模型没有限制,但其在Java应用中的频繁使用让它备受青睐,以至于后来让它作为EJB(EnterpriseJavaBeans)模型的补充,甚至是替补。Spring框架为开发提供了一系列的解决方案,比如利用控制反转的核心特性,并通过依赖注入实现控制反转来实现管理对象生命周期容器化,利用面向切面编程进行声明式的事务管理,整合多种持久化技术管理数据访问,提供大量优秀的Web框架方便开发等等。Spring框架具有控制反转(IOC)特性,IOC旨在方便项目维护和测试,它提供了一种通过Java的反射机制对Java对象进行统一的配置和管理的方法。Spring框架利用容器管理对象的生命周期,容器可以通过扫描XML文件或类上特定Java注解来配置对象,开发者可以通过依赖查找或依赖注入来获得对象。Spring框架具有面向切面编程(AOP)框架,SpringAOP框架基于代理模式,同时运行时可配置;AOP框架主要针对模块之间的交叉关注点进行模块化。Spring框架的AOP框架仅提供基本的AOP特性,虽无法与AspectJ框架相比,但通过与AspectJ的集成,也可以满足基本需求。Spring框架下的事务管理、远程访问等功能均可以通过使用SpringAOP技术实现。Spring的事务管理框架为Java平台带来了一种抽象机制,使本地和全局事务以及嵌套事务能够与保存点一起工作,并且几乎可以在Java平台的任何环境中工作。Spring集成多种事务模板,系统可以通过事务模板、XML或Java注解进行事务配置,并且事务框架集成了消息传递和缓存等功能。Spring的数据访问框架解决了开发人员在应用程序中使用数据库时遇到的常见困难。它不仅对Java:JDBC、iBATS/MyBATIs、Hibernate、Java数据对象(JDO)、ApacheOJB和ApacheCayne等所有流行的数据访问框架中提供支持,同时还可以与Spring的事务管理一起使用,为数据访问提供了灵活的抽象。Spring框架最初是没有打算构建一个自己的WebMVC框架,其开发人员在开发过程中认为现有的StrutsWeb框架的呈现层和请求处理层之间以及请求处理层和模型之间的分离不够,于是创建了SpringMVC。SpringBoot整合kafka 消费者 a. 引入Pom b.JAVA代码 配置文件(YML) 消费者 a.Pom文件 b1.Java代码(无并发访问、无批量获取) b2.配置文件 c.Java代码(并发、批量获取) Kafka安装 Kafka是由Apache软件基金会开发的一个开源流处理平台,是一种高吞吐量的分布式发布订阅消息系统。 主要包含几个组件: Topic:消息主题,特定消息的发布接口,每个Topic都可以分成数个Partition,用于消息的并发发送。 Producer:生产者,信息的发布者,发布者可以指定数个Partition进行发布。 Consumer:消费者,信息的使用者,同一个Group的消费者数量,最好不好超过Partition的数量,对于分区的Topic,消费者使用时需要指定相应的分区号。 Broker:服务代理 ##下载kafka SpringBoot整合kafka 当前SpringBoot版本为2.0.2.RELEASE,打包工具为Maven 消费者 a. 引入Pom <?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.kafkatest</groupId> <artifactId>producer</artifactId> <version>1.0-SNAPSHOT</version> <name>kafka-producer</name> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.0.2.RELEASE</version> <relativePath/> </parent> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> <java.version>1.8</java.version> <joda-time.version>2.3</joda-time.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> </project> b.JAVA代码 @Service public class KafkaProducerTest { @Autowired private KafkaTemplate<String,byte[]> kafkaTemplate; private final String topic = "byteArray_topic1"; public void sendMessage(int key,String value){ ProducerRecord<String,byte[]> record = new ProducerRecord<>(topic, key%3,String.valueOf(key),value.getBytes()); kafkaTemplate.send(record); } } 配置文件(YML) spring: kafka: producer: bootstrap-servers: 172.169.0.109:9092 batch-size: 16384 retries: 0 buffer-memory: 33554432 key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.ByteArraySerializer 这里有一个非常陷阱的问题需要特别注意:序列化类的路径是:org.apache.kafka.common.serialization.StringSerializer 而不是 org.apache.kafka.config.serialization.StringSerializer 否则会出现如下错误: 2019-01-31 11:35:14.794 [main] WARN o.s.c.a.AnnotationConfigApplicationContext - Exception encountered during context initialization - cancelling refresh attempt: org.springframework.beans.factory.UnsatisfiedDependencyException: Error creating bean with name 'kafkaProducerTest': Unsatisfied dependency expressed through field 'kafkaTemplate'; nested exception is org.springframework.beans.factory.UnsatisfiedDependencyException: Error creating bean with name 'org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration': Unsatisfied dependency expressed through constructor parameter 0; nested exception is org.springframework.boot.context.properties.ConfigurationPropertiesBindException: Error creating bean with name 'spring.kafka-org.springframework.boot.autoconfigure.kafka.KafkaProperties': Could not bind properties to 'KafkaProperties' : prefix=spring.kafka, ignoreInvalidFields=false, ignoreUnknownFields=true; nested exception is org.springframework.boot.context.properties.bind.BindException: Failed to bind properties under 'spring.kafka.producer.key-serializer' to java.lang.Class<?> 2019-01-31 11:35:14.810 [main] ERROR o.s.b.d.LoggingFailureAnalysisReporter - *************************** APPLICATION FAILED TO START *************************** Description: Failed to bind properties under 'spring.kafka.producer.key-serializer' to java.lang.Class<?>: Property: spring.kafka.producer.key-serializer Value: org.apache.kafka.config.serialization.StringSerializer Origin: class path resource [application.yml]:8:25 Reason: No converter found capable of converting from type [java.lang.String] to type [java.lang.Class<?>] Action: Update your application's configuration 消费者 如果不使用并发获取、批量获取消费者的代码非常简单。 a.Pom文件 <?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.kafkatest</groupId> <artifactId>consumer</artifactId> <version>1.0-SNAPSHOT</version> <name>kafka-consumer</name> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.0.2.RELEASE</version> <relativePath/> <!-- lookup parent from repository --> </parent> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> <java.version>1.8</java.version> <joda-time.version>2.3</joda-time.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> </project> b1.Java代码(无并发访问、无批量获取) @Service @Slf4j public class Listener { private final String topic = "byteArray_topic1"; public void listen(ConsumerRecord<String, byte[]> record){ log.info("kafka的key: " + record.key()); log.info("kafka的value: " + new String(record.value())); } } b2.配置文件 spring: kafka: consumer: enable-auto-commit: true group-id: gridMonitorGroup auto-commit-interval: 1000 auto-offset-reset: latest bootstrap-servers: "172.169.0.109:9092" key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.ByteArrayDeserializer c.Java代码(并发、批量获取) Kafka消费者配置类 批量获取关键代码: ①factory.setBatchListener(true); ②propsMap.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG,50); 并发获取关键代码: factory.setConcurrency(concurrency); @Configuration @EnableKafka public class KafkaConsumerConfig { @Value("${kafka.consumer.bootstrap-servers}") private String servers; @Value("${kafka.consumer.enable-auto-commit}") private boolean enableAutoCommit; @Value("${kafka.consumer.auto-commit-interval}") private String autoCommitInterval; @Value("${kafka.consumer.group-id}") private String groupId; @Value("${kafka.consumer.auto-offset-reset}") private String autoOffsetReset; @Value("${kafka.consumer.concurrency}") private int concurrency; @Bean public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, byte[]>> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, byte[]> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); //并发数量 factory.setConcurrency(concurrency); //批量获取 factory.setBatchListener(true); factory.getContainerProperties().setPollTimeout(1500); return factory; } public ConsumerFactory<String, byte[]> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfigs()); } public Map<String, Object> consumerConfigs() { Map<String, Object> propsMap = new HashMap<>(); propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers); propsMap.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, enableAutoCommit); propsMap.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, autoCommitInterval); propsMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); propsMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); propsMap.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); propsMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset); //最多批量获取50个 propsMap.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG,50); return propsMap; } @Bean public Listener listener() { return new Listener(); } } Kafka消费者Listener @Service @Slf4j public class Listener { private final String topic = "byteArray_topic1"; @KafkaListener(id="myListener", topicPartitions ={@TopicPartition(topic = topic, partitions = { "0", "1" ,"2"})}) public void listen(List<ConsumerRecord<String, byte[]>> recordList) { recordList.forEach((record)->{ log.info("kafka的key: " + record.key()); log.info("kafka的value: " + new String(record.value())); }); } } 配置文件 kafka: consumer: enable-auto-commit: true group-id: gridMonitorGroup auto-commit-interval: 1000 auto-offset-reset: latest bootstrap-servers: "172.169.0.109:9092" key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.ByteArrayDeserializer concurrency: 3 ———————————————— 原文链接:https://blog.csdn.net/menxin_job/article/details/86712973
-
本篇内容主要讲解“Springboot如何集成Kafka进行批量消费”,感兴趣的朋友不妨来看看。本文介绍的方法操作简单快捷,实用性强。下面就让小编来带大家学习“Springboot如何集成Kafka进行批量消费”吧!Spring Boot是由Pivotal团队提供的全新框架,其设计目的是用来简化新Spring应用的初始搭建以及开发过程。该框架使用了特定的方式来进行配置,从而使开发人员不再需要定义样板化的配置。通过这种方式,Spring Boot致力于在蓬勃发展的快速应用开发领域(rapid application development)成为领导者。Spring框架是Java平台上的一种开源应用框架,提供具有控制反转特性的容器。尽管Spring框架自身对编程模型没有限制,但其在Java应用中的频繁使用让它备受青睐,以至于后来让它作为EJB(EnterpriseJavaBeans)模型的补充,甚至是替补。Spring框架为开发提供了一系列的解决方案,比如利用控制反转的核心特性,并通过依赖注入实现控制反转来实现管理对象生命周期容器化,利用面向切面编程进行声明式的事务管理,整合多种持久化技术管理数据访问,提供大量优秀的Web框架方便开发等等。Spring框架具有控制反转(IOC)特性,IOC旨在方便项目维护和测试,它提供了一种通过Java的反射机制对Java对象进行统一的配置和管理的方法。Spring框架利用容器管理对象的生命周期,容器可以通过扫描XML文件或类上特定Java注解来配置对象,开发者可以通过依赖查找或依赖注入来获得对象。Spring框架具有面向切面编程(AOP)框架,SpringAOP框架基于代理模式,同时运行时可配置;AOP框架主要针对模块之间的交叉关注点进行模块化。Spring框架的AOP框架仅提供基本的AOP特性,虽无法与AspectJ框架相比,但通过与AspectJ的集成,也可以满足基本需求。Spring框架下的事务管理、远程访问等功能均可以通过使用SpringAOP技术实现。Spring的事务管理框架为Java平台带来了一种抽象机制,使本地和全局事务以及嵌套事务能够与保存点一起工作,并且几乎可以在Java平台的任何环境中工作。Spring集成多种事务模板,系统可以通过事务模板、XML或Java注解进行事务配置,并且事务框架集成了消息传递和缓存等功能。Spring的数据访问框架解决了开发人员在应用程序中使用数据库时遇到的常见困难。它不仅对Java:JDBC、iBATS/MyBATIs、Hibernate、Java数据对象(JDO)、ApacheOJB和ApacheCayne等所有流行的数据访问框架中提供支持,同时还可以与Spring的事务管理一起使用,为数据访问提供了灵活的抽象。Spring框架最初是没有打算构建一个自己的WebMVC框架,其开发人员在开发过程中认为现有的StrutsWeb框架的呈现层和请求处理层之间以及请求处理层和模型之间的分离不够,于是创建了SpringMVC。引入依赖<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> <version>1.3.11.RELEASE</version> </dependency>因为我的项目的 springboot 版本是 1.5.22.RELEASE,所以引的是 1.3.11.RELEASE 的包。读者可以根据下图来自行选择对应的版本。图片更新可能不及时,详情可查看spring-kafka 官方网站。注:这里有个踩坑点,如果引入包版本不对,项目启动时会抛出org.springframework.core.log.LogAccessor 异常:java.lang.ClassNotFoundException: org.springframework.core.log.LogAccessor创建配置类/** * kafka 配置类 */ @Configuration @EnableKafka public class KafkaConsumerConfig { private static final org.slf4j.Logger LOGGER = LoggerFactory.getLogger(KafkaConsumerConfig.class); @Value("${kafka.bootstrap.servers}") private String kafkaBootstrapServers; @Value("${kafka.group.id}") private String kafkaGroupId; @Value("${kafka.topic}") private String kafkaTopic; public static final String CONFIG_PATH = "/home/admin/xxx/BOOT-INF/classes/kafka_client_jaas.conf"; public static final String LOCATION_PATH = "/home/admin/xxx/BOOT-INF/classes/kafka.client.truststore.jks"; @Bean public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); // 设置并发量,小于或者等于 Topic 的分区数 factory.setConcurrency(5); // 设置为批量监听 factory.setBatchListener(Boolean.TRUE); factory.getContainerProperties().setPollTimeout(30000); return factory; } public ConsumerFactory<String, String> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfigs()); } public Map<String, Object> consumerConfigs() { Map<String, Object> props = new HashMap<>(); //设置接入点,请通过控制台获取对应Topic的接入点。 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBootstrapServers); //设置SSL根证书的路径,请记得将XXX修改为自己的路径。 //与SASL路径类似,该文件也不能被打包到jar中。 System.setProperty("java.security.auth.login.config", CONFIG_PATH); props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, LOCATION_PATH); //根证书存储的密码,保持不变。 props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient"); //接入协议,目前支持使用SASL_SSL协议接入。 props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL"); //SASL鉴权方式,保持不变。 props.put(SaslConfigs.SASL_MECHANISM, "PLAIN"); // 自动提交 props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, Boolean.TRUE); //两次Poll之间的最大允许间隔。 //消费者超过该值没有返回心跳,服务端判断消费者处于非存活状态,服务端将消费者从Consumer Group移除并触发Rebalance,默认30s。 props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000); //设置单次拉取的量,走公网访问时,该参数会有较大影响。 props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 32000); props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, 32000); //每次Poll的最大数量。 //注意该值不要改得太大,如果Poll太多数据,而不能在下次Poll之前消费完,则会触发一次负载均衡,产生卡顿。 props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 30); //消息的反序列化方式。 props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); //当前消费实例所属的消费组,请在控制台申请之后填写。 //属于同一个组的消费实例,会负载消费消息。 props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaGroupId); //Hostname校验改成空。 props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, ""); return props; } }注:此处通过 factory.setConcurrency(5); 配置了并发量为 5 ,假设我们线上的 Topic 有 12 个分区。那么将会是 3 个线程分配到 2 个分区,2 个线程分配到 3 个分区,3 * 2 + 2 * 3 = 12。Kafka 消费者/** * kafka 消息消费类 */ @Component public class KafkaMessageListener { private static final Logger LOGGER = LoggerFactory.getLogger(KafkaMessageListener.class); @KafkaListener(topics = {"${kafka.topic}"}) public void listen(List<ConsumerRecord<String, String>> recordList) { for (ConsumerRecord<String,String> record : recordList) { // 打印消息的分区以及偏移量 LOGGER.info("Kafka Consume partition:{}, offset:{}", record.partition(), record.offset()); String value = record.value(); System.out.println("value = " + value); // 处理业务逻辑 ... } } }因为我在配置类中设置了批量监听,所以此处 listen 方法的入参是List:List<ConsumerRecord<String, String>>。到此,相信大家对“Springboot如何集成Kafka进行批量消费”有了更深的了解,不妨来实际操作一番吧!原文链接:https://www.yisu.com/zixun/622113.html
-
参数校验失败会自动引发异常,我们当然不可能再去手动捕捉异常进行处理。但我们又不想手动捕捉这个异常,又要对这个异常进行处理,那正好使用SpringBoot全局异常处理来达到一劳永逸的效果!1、基本使用首先,我们需要新建一个类,在这个类上加上@ControllerAdvice或@RestControllerAdvice注解,这个类就配置成全局处理类了。这个根据你的Controller层用的是@Controller还是@RestController来决定。然后在类中新建方法,在方法上加上@ExceptionHandler注解并指定你想处理的异常类型,接着在方法内编写对该异常的操作逻辑,就完成了对该异常的全局处理!我们现在就来演示一下对参数校验失败抛出的MethodArgumentNotValidException全局处理:package com.csdn.demo1.global; import org.springframework.validation.ObjectError; import org.springframework.web.bind.MethodArgumentNotValidException; import org.springframework.web.bind.annotation.ExceptionHandler; import org.springframework.web.bind.annotation.RestControllerAdvice; @RestControllerAdvice @ResponseBody public class ExceptionControllerAdvice { @ExceptionHandler(MethodArgumentNotValidException.class) @ResponseStatus(HttpStatus.BAD_REQUEST) public String MethodArgumentNotValidExceptionHandler(MethodArgumentNotValidException e) { // 从异常对象中拿到ObjectError对象 ObjectError objectError = e.getBindingResult().getAllErrors().get(0); // 然后提取错误提示信息进行返回 return objectError.getDefaultMessage(); } /** * 系统异常 预期以外异常 */ @ExceptionHandler(Exception.class) @ResponseStatus(value = HttpStatus.INTERNAL_SERVER_ERROR) public ResultVO<?> handleUnexpectedServer(Exception ex) { log.error("系统异常:", ex); // GlobalMsgEnum.ERROR是我自己定义的枚举类 return new ResultVO<>(GlobalMsgEnum.ERROR); } /** * 所以异常的拦截 */ @ExceptionHandler(Throwable.class) @ResponseStatus(value = HttpStatus.INTERNAL_SERVER_ERROR) public ResultVO<?> exception(Throwable ex) { log.error("系统异常:", ex); return new ResultVO<>(GlobalMsgEnum.ERROR); } }我们再次进行测试,这次返回的就是我们制定的错误提示信息!我们通过全局异常处理优雅的实现了我们想要的功能!以后我们再想写接口参数校验,就只需要在入参的成员变量上加上Validator校验规则注解,然后在参数上加上@Valid注解即可完成校验,校验失败会自动返回错误提示信息,无需任何其他代码!2、自定义异常在很多情况下,我们需要手动抛出异常,比如在业务层当有些条件并不符合业务逻辑,而使用自定义异常有诸多优点:自定义异常可以携带更多的信息,不像这样只能携带一个字符串。项目开发中经常是很多人负责不同的模块,使用自定义异常可以统一了对外异常展示的方式。自定义异常语义更加清晰明了,一看就知道是项目中手动抛出的异常。我们现在就来开始写一个自定义异常:package com.csdn.demo1.global; import lombok.Getter; @Getter //只要getter方法,无需setter public class APIException extends RuntimeException { private int code; private String msg; public APIException() { this(1001, "接口错误"); } public APIException(String msg) { this(1001, msg); } public APIException(int code, String msg) { super(msg); this.code = code; this.msg = msg; } }然后在刚才的全局异常类中加入如下://自定义的全局异常 @ExceptionHandler(APIException.class) public String APIExceptionHandler(APIException e) { return e.getMsg(); }这样就对异常的处理就比较规范了,当然还可以添加对Exception的处理,这样无论发生什么异常我们都能屏蔽掉然后响应数据给前端,不过建议最后项目上线时这样做,能够屏蔽掉错误信息暴露给前端,在开发中为了方便调试还是不要这样做。另外,当我们抛出自定义异常的时候全局异常处理只响应了异常中的错误信息msg给前端,并没有将错误代码code返回。这还需要配合数据统一响应。如果在多模块使用,全局异常等公共功能抽象成子模块,则在需要的子模块中需要将该模块包扫描加入,@SpringBootApplication(scanBasePackages = {"com.xxx"})
-
1、介绍一个接口一般对参数(请求数据)都会进行安全校验,参数校验的重要性自然不必多说,那么如何对参数进行校验就有讲究了。一般来说有三种常见的校验方式,我们使用了最简洁的第三种方法业务层校验Validator + BindResult校验Validator + 自动抛出异常业务层校验无需多说,即手动在java的Service层进行数据校验判断。不过这样太繁琐了,光校验代码就会有很多而使用Validator+ BindingResult已经是非常方便实用的参数校验方式了,在实际开发中也有很多项目就是这么做的,不过这样还是不太方便,因为你每写一个接口都要添加一个BindingResult参数,然后再提取错误信息返回给前端(简单看一下)。@PostMapping("/addUser") public String addUser(@RequestBody @Validated User user, BindingResult bindingResult) { // 如果有参数校验失败,会将错误信息封装成对象组装在BindingResult里 List<ObjectError> allErrors = bindingResult.getAllErrors(); if(!allErrors.isEmpty()){ return allErrors.stream() .map(o->o.getDefaultMessage()) .collect(Collectors.toList()).toString(); } // 返回默认的错误信息 // return allErrors.get(0).getDefaultMessage(); return validationService.addUser(user); }2、Validator + 自动抛出异常(使用)内置参数校验如下:首先Validator可以非常方便的制定校验规则,并自动帮你完成校验。首先在入参里需要校验的字段加上注解,每个注解对应不同的校验规则,并可制定校验失败后的信息:@Data public class User { @NotNull(message = "用户id不能为空") private Long id; @NotNull(message = "用户账号不能为空") @Size(min = 6, max = 11, message = "账号长度必须是6-11个字符") private String account; @NotNull(message = "用户密码不能为空") @Size(min = 6, max = 11, message = "密码长度必须是6-16个字符") private String password; @NotNull(message = "用户邮箱不能为空") @Email(message = "邮箱格式不正确") private String email; }校验规则和错误提示信息配置完毕后,接下来只需要在接口仅需要在校验的参数上加上@Valid注解(去掉BindingResult后会自动引发异常,异常发生了自然而然就不会执行业务逻辑):@RestController @RequestMapping("user") public class ValidationController { @Autowired private ValidationService validationService; @PostMapping("/addUser") public String addUser(@RequestBody @Validated User user) { return validationService.addUser(user); } }现在我们进行测试,打开knife4j文档地址,当输入的请求数据为空时,Validator会将所有的报错信息全部进行返回,所以需要与全局异常处理一起使用。// 使用form data方式调用接口,校验异常抛出 BindException // 使用 json 请求体调用接口,校验异常抛出 MethodArgumentNotValidException // 单个参数校验异常抛出ConstraintViolationException // 处理 json 请求体调用接口校验失败抛出的异常 @ExceptionHandler(MethodArgumentNotValidException.class) public ResultVO<String> MethodArgumentNotValidException(MethodArgumentNotValidException e) { List<FieldError> fieldErrors = e.getBindingResult().getFieldErrors(); List<String> collect = fieldErrors.stream() .map(DefaultMessageSourceResolvable::getDefaultMessage) .collect(Collectors.toList()); return new ResultVO(ResultCode.VALIDATE_FAILED, collect); } // 使用form data方式调用接口,校验异常抛出 BindException @ExceptionHandler(BindException.class) public ResultVO<String> BindException(BindException e) { List<FieldError> fieldErrors = e.getBindingResult().getFieldErrors(); List<String> collect = fieldErrors.stream() .map(DefaultMessageSourceResolvable::getDefaultMessage) .collect(Collectors.toList()); return new ResultVO(ResultCode.VALIDATE_FAILED, collect); }3、分组校验和递归校验分组校验有三个步骤:定义一个分组类(或接口)在校验注解上添加groups属性指定分组Controller方法的@Validated注解添加分组类public interface Update extends Default{ }@Data public class User { @NotNull(message = "用户id不能为空",groups = Update.class) private Long id; ...... }@PostMapping("update") public String update(@Validated({Update.class}) User user) { return "success"; }如果Update不继承Default,@Validated({Update.class})就只会校验属于Update.class分组的参数字段;如果继承了,会校验了其他默认属于Default.class分组的字段。对于递归校验(比如类中类),只要在相应属性类上增加@Valid注解即可实现(对于集合同样适用)4、自定义校验Spring Validation允许用户自定义校验,实现很简单,分两步:自定义校验注解编写校验者类@Target({ METHOD, FIELD, ANNOTATION_TYPE, CONSTRUCTOR, PARAMETER, TYPE_USE }) @Retention(RUNTIME) @Documented @Constraint(validatedBy = {HaveNoBlankValidator.class})// 标明由哪个类执行校验逻辑 public @interface HaveNoBlank { // 校验出错时默认返回的消息 String message() default "字符串中不能含有空格"; Class<?>[] groups() default { }; Class<? extends Payload>[] payload() default { }; /** * 同一个元素上指定多个该注解时使用 */ @Target({ METHOD, FIELD, ANNOTATION_TYPE, CONSTRUCTOR, PARAMETER, TYPE_USE }) @Retention(RUNTIME) @Documented public @interface List { NotBlank[] value(); } }public class HaveNoBlankValidator implements ConstraintValidator<HaveNoBlank, String> { @Override public boolean isValid(String value, ConstraintValidatorContext context) { // null 不做检验 if (value == null) { return true; } // 校验失败 return !value.contains(" "); // 校验成功 } }
-
Java本地缓存 Java实现本地缓存的方式有很多,其中比较常见的有HashMap、Guava Cache、Caffeine和Encahche等。这些缓存技术各有优缺点,你可以根据自己的需求选择适合自己的缓存技术。以下是一些详细介绍: 1. HashMap:通过Map的底层方式,直接将需要缓存的对象放在内存中。优点是简单粗暴,不需要引入第三方包,比较适合一些比较简单的场景。缺点是没有缓存淘汰策略,定制化开发成本高。 2. Guava Cache:Guava是一个Google开源的项目,提供了一些Java工具类和库。Guava Cache是Guava提供的一个本地缓存框架,它使用LRU算法来管理缓存。优点是性能好,支持异步加载和批量操作。缺点是需要引入Guava库。 3. Caffeine:Caffeine是一个高性能的Java本地缓存库,它使用了基于时间戳的过期策略和可扩展性设计。优点是性能好,支持异步加载和批量操作。缺点是需要引入Caffeine库。 4. Encahche:Encahche是一个轻量级的Java本地缓存库,它使用了基于时间戳的过期策略和可扩展性设计。优点是性能好,支持异步加载和批量操作。缺点是需要引入Encahche库。 示例代码 1.Guava Cache示例代码 以下是使用Guava Cache实现Java本地缓存的示例代码: import com.google.common.cache.CacheBuilder; import com.google.common.cache.CacheLoader; import com.google.common.cache.LoadingCache; import java.util.concurrent.ExecutionException; public class GuavaCacheExample { private static final LoadingCache<String, String> CACHE = CacheBuilder.newBuilder() .maximumSize(100) // 设置缓存最大容量为100 .build(new CacheLoader<String, String>() { @Override public String load(String key) throws Exception { // 从数据库中查询数据并返回 return queryDataFromDatabase(key); } }); public static void main(String[] args) throws Exception { // 从缓存中获取数据 String data = CACHE.get("key"); System.out.println(data); } } 2.Caffeine示例代码 以下是使用Caffeine实现Java本地缓存的示例代码: import com.github.benmanes.caffeine.cache.Cache; import com.github.benmanes.caffeine.cache.Caffeine; import java.util.concurrent.TimeUnit; public class CaffeineExample { private static final Cache<String, String> CACHE = Caffeine.newBuilder() .maximumSize(100) // 设置缓存最大容量为100 .expireAfterWrite(10, TimeUnit.MINUTES) // 设置缓存过期时间为10分钟 .build(); public static void main(String[] args) throws Exception { // 从缓存中获取数据 String data = CACHE.get("key", new Callable<String>() { @Override public String call() throws Exception { // 从数据库中查询数据并返回 return queryDataFromDatabase("key"); } }); System.out.println(data); } } 3.Encahche示例代码 以下是使用Encahche实现Java本地缓存的示例代码: import org.ehcache.Cache; import org.ehcache.CacheManager; import org.ehcache.config.Builder; import org.ehcache.config.Configuration; import org.ehcache.config.units.MemoryUnit; public class EncahcheExample { private static final Cache<String, String> CACHE = CacheManager.create() .newCache("myCache", new Configuration() .withSizeOfMaxObjectSize(1024 * 1024) // 设置缓存最大容量为1MB .withExpiry(10, TimeUnit.MINUTES)) // 设置缓存过期时间为10分钟 .build(); public static void main(String[] args) throws Exception { // 从缓存中获取数据 String data = CACHE.get("key"); System.out.println(data); } }
-
Java线程池是一种管理线程的机制,它可以有效地控制并发执行的线程数量,提高程序的性能和稳定性。本文将介绍Java线程池的概念、实现原理以及一个简单的示例代码。一、Java线程池概念线程池的作用:线程池可以预先创建一定数量的线程,当有任务需要执行时,从线程池中获取一个空闲的线程来执行任务,任务执行完毕后,将线程归还给线程池。这样可以避免频繁地创建和销毁线程,提高系统的性能。线程池的优点:提高系统性能:通过复用线程,减少了线程创建和销毁的开销。控制并发数:通过限制线程池中的线程数量,可以有效地控制并发执行的线程数量,避免过多的线程导致系统资源耗尽。提高响应速度:由于线程池中的线程已经预先创建好,所以在接收到任务时可以直接使用,无需等待线程创建完成。二、Java线程池实现原理Java线程池主要由以下几个部分组成:核心线程数(corePoolSize):线程池中始终保持的线程数量。当任务队列已满且有新任务到来时,会创建新的线程来执行任务,直到达到核心线程数上限时,后续的任务将被放入阻塞队列中等待执行。最大线程数(maximumPoolSize):线程池允许的最大线程数量。当阻塞队列已满且有新任务到来时,如果核心线程数已达到最大值,那么新任务将会被放入拒绝队列中。工作队列(workQueue):用于存放等待执行的任务的阻塞队列。常用的阻塞队列有ArrayBlockingQueue、LinkedBlockingQueue等。拒绝策略(rejectedExecutionHandler):当阻塞队列已满且无法创建新线程时,如何处理新任务的方法。常用的拒绝策略有AbortPolicy(抛出异常)、DiscardPolicy(丢弃任务)和CallerRunsPolicy(让调用者自己执行任务)。线程工厂(threadFactory):用于创建新线程的工厂类。可以通过自定义线程工厂为每个新创建的线程设置一些属性,如名称、优先级等。三、Java线程池示例代码Executors.newFixedThreadPool:创建一个固定大小的线程池,可控制并发的线程数,超出的线程会在队列中等待。 Executors.newCachedThreadPool:创建一个可缓存的线程池,若线程数超过处理所需,缓存一段时间后会回收,若线程数不够,则新建线程。 Executors.newSingleThreadExecutor:创建单个线程数的线程池,它可以保证先进先出的执行顺序。 Executors.newScheduledThreadPool:创建一个可以执行延迟任务的线程池。 Executors.newSingleThreadScheduledExecutor:创建一个单线程的可以执行延迟任务的线程池。 Executors.newWorkStealingPool:创建一个抢占式执行的线程池(任务执行顺序不确定)【JDK 1.8 添加】。 ThreadPoolExecutor:手动创建线程池的方式,它创建时最多可以设置 7 个参数。 下面是Java线程池示例代码: 以下是各种创建线程池的示例代码: 以下是各种创建线程池的示例代码:创建一个固定大小的线程池:ExecutorService fixedThreadPool = Executors.newFixedThreadPool(5); fixedThreadPool.execute(() -> System.out.println("Hello, world!")); fixedThreadPool.shutdown(); // 必须调用 shutdown() 才能关闭线程池。否则会导致内存泄漏。 创建一个可缓存的线程池:ExecutorService cachedThreadPool = Executors.newCachedThreadPool(); cachedThreadPool.execute(() -> System.out.println("Hello, world!")); cachedThreadPool.shutdown(); // 必须调用 shutdown() 才能关闭线程池。否则会导致内存泄漏。 创建一个单个线程的线程池:ExecutorService singleThreadExecutor = Executors.newSingleThreadExecutor(); singleThreadExecutor.execute(() -> System.out.println("Hello, world!")); singleThreadExecutor.shutdown(); // 必须调用 shutdown() 才能关闭线程池。否则会导致内存泄漏。 创建一个可以执行延迟任务的线程池:ScheduledExecutorService scheduledThreadPool = Executors.newScheduledThreadPool(2); scheduledThreadPool.scheduleAtFixedRate(() -> System.out.println("Hello, world!"), 0, 5, TimeUnit.SECONDS); scheduledThreadPool.shutdown(); // 必须调用 shutdown() 才能关闭线程池。否则会导致内存泄漏。 创建一个单线程的可以执行延迟任务的线程池:ScheduledExecutorService singleThreadScheduledExecutor = Executors.newSingleThreadScheduledExecutor(); singleThreadScheduledExecutor.scheduleAtFixedRate(() -> System.out.println("Hello, world!"), 0, 5, TimeUnit.SECONDS); singleThreadScheduledExecutor.shutdown(); // 必须调用 shutdown() 才能关闭线程池。否则会导致内存泄漏。
-
当您看到此通知时,它表示 CodeArts IDE代码文件观察器的句柄已用完,因为工作区很大并且包含许多文件。在调整平台限制之前,请确保将可能较大的文件夹(例如 Python .venv)添加到files.watcherExclude设置中(更多详细信息见下文)。可以通过运行查看当前限制:cat /proc/sys/fs/inotify/max_user_watches可以通过编辑/etc/sysctl.conf(Arch Linux 除外,阅读下文)并将如下行添加到文件末尾来将限制增加到最大值:fs.inotify.max_user_watches=524288然后可以通过运行加载新值sudo sysctl -p。虽然 524,288 是可以观看的最大文件数,但如果您处于内存特别受限的环境中,则可能需要降低该数字。每个文件 watch占用 1080 个字节,因此假设所有 524,288 个 watch 都被消耗,结果上限约为 540 MiB。基于Arch的发行版(包括 Manjaro)要求您更改不同的文件;请按照以下步骤操作。files.watcherExclude 另一种选择是使用设置从CodeArts IDE代码文件观察器中排除特定的工作区目录。默认为files.watcherExcludeexcludesnode_modules和 下的一些文件夹.git,但您可以添加其他不希望 CodeArts IDE跟踪的目录。"files.watcherExclude": { "**/.git/objects/**": true, "**/.git/subtree-cache/**": true, "**/node_modules/*/**": true }
-
Java8使用stream流给List<Map<String,Object>>根据字段key分组 一、项目场景: 从已得到的List集合中,根据某一元素(这里指map的key)进行分组,筛选出需要的数据。 如果是SQL的话则使用group by直接实现,代码的方式则如下: 使用到stream流的Collectors.groupingBy()方法。 二、代码实现 1、首先将数据add封装到List中,完成数据准备。 //groupList用于库-表分组的list,减少jdbc连接时间 List<Map<String,Object>> groupList = new ArrayList<>(); Map<String,Object> map1 = new HashMap<>(); map.put("name","张三"); map.put("age",20); Map<String,Object> map2 = new HashMap<>(); map.put("name","李四"); map.put("age",20); //excel每行的值存入集合中 groupList.add(map1); groupList.add(map2); 2、然后按照name属性进行分组(单字段),使用stream流分组 //分组判断,stream流 Map<String, List<Map<String, Object>>> listMap = groupList.stream().collect( Collectors.groupingBy(item -> item.get("name").toString()) ); ··按照name,age属性进行分组(多字段),使用stream流分组 //分组判断,stream流 Map<String, List<Map<String, Object>>> listMap = groupList.stream().collect( Collectors.groupingBy(item -> item.get("name").toString()+"|"+item.get("age")) ); 3、遍历结果,输出查看 得到Map<String, List<Map<String, Object>>>, key是你分组的字段,value分组下对应的值。 也就是group by的效果。 for (String groupKey : listMap.keySet()) { //分组key System.out.println("分组Key: "+groupKey); System.out.println("分组Key的value: "+listMap.get("groupKey")); } ———————————————— 版权声明:本文为CSDN博主「Lin'ZT」的原创文章,遵循CC 4.0 BY-SA版权协议,转载请附上原文出处链接及本声明。 原文链接:https://blog.csdn.net/weixin_43920527/article/details/129820467
-
默认情况下,Spring的RestTemplate发送POST请求时,请求Content-Type类型是application/json。如果要将请求的Content-Type设置为form-data,可以使用MultiValueMap,并使用postForObject方法发送POST请求。示例如下:import org.springframework.http.HttpEntity; import org.springframework.http.HttpHeaders; import org.springframework.http.MediaType; import org.springframework.util.LinkedMultiValueMap; import org.springframework.util.MultiValueMap; import org.springframework.web.client.RestTemplate; public class RestClient { private RestTemplate restTemplate; public RestClient() { this.restTemplate = new RestTemplate(); } public void doPostRequest() { HttpHeaders headers = new HttpHeaders(); headers.setContentType(MediaType.MULTIPART_FORM_DATA); MultiValueMap<String, String> body = new LinkedMultiValueMap<>(); body.add("key1", "value1"); body.add("key2", "value2"); HttpEntity<MultiValueMap<String, String>> requestEntity = new HttpEntity<>(body, headers); String url = "http://your-api-url"; String response = restTemplate.postForObject(url, requestEntity, String.class); System.out.println(response); } }复制在上面的代码中,我们首先创建了一个HttpHeaders对象,将ContentType设置为MediaType.MULTIPART_FORM_DATA,表示该请求是form-data类型。然后,我们使用LinkedMultiValueMap将请求的参数添加到请求体中。最后,我们使用postForObject方法向API发送POST请求,并传入url、requestEntity和返回值类型。MultiValueMap是一个类似于Map的数据结构,它在键-值对中允许存在多个值。在Java语言中,Spring框架中的org.springframework.util.MultiValueMap接口定义了一个具有此特性的映射表数据类型。对于普通的Map接口来说,在一个键(key)下只能包含一个值(value),而MultiValueMap可以容纳多个值,例如:MultiValueMap<String, String> multiValueMap = new LinkedMultiValueMap<>(); multiValueMap.add("name", "张三"); multiValueMap.add("hobby", "篮球"); multiValueMap.add("hobby", "足球"); System.out.println(multiValueMap);复制上面代码中我们使用LinkedMultiValueMap创建了一个MultiValueMap类型的对象,并向其添加了键为"name"和"hobby"的数据,其中"hobby"这个键下有两个值。输出结果为:{name=[张三], hobby=[篮球, 足球]}复制该数据结构主要用于解析HTML表单提交的数据或HTTP请求参数,特别是当前端页面的某个表单字段需要提交一个列表或数组时。使用MultiValueMap则无需手动处理重复键值,可以更方便地进行数据管理和交换。Spring框架中的RestTemplate类,就利用了MultiValueMap实现了HTTP请求参数的传递和解析。其他一些框架和库,如Jackson、Gson等也支持将JSON转换为MultiValueMap,并方便地将其转换为Java对象。
-
RestTemplate 和 HttpClient 都是Java中用于发送HTTP请求的工具之一,但它们在实现上有些不同。HttpClient 是 Apache 开源组织提供的一个URL客户端库。它提供了HTTP协议的客户端实现。可以通过 HttpPost,HttpGet,HttpPut等类实现相应的 HTTP 请求,同时 HttpClient 提供了 Header,Response等信息的处理。使用 HttpClient 可以控制HTTP请求与响应行为的各个方面,比如连接超时,重定向,异常处理等。而 RestTemplate 是 Spring Framework 提供的模板类,用于简化 HTTP 请求与响应处理。它是基于 HttpMessageConverter 实现的,可以将 HTTP 消息转换为 Java 对象和将 Java 对象转换为 HTTP 消息。它支持 GET,POST,PUT,DELETE,PATCH 等 HTTP 方法。使用 RestTemplate 可以更加方便地进行 HTTP 请求,同时它也提供了自定义的拦截器、错误处理和重试机制等功能。相比较而言,RestTemplate 更适合于集成Spring框架的项目中使用,其简化代码结构,生产力更高,可进行更丰富的配置和扩展;而 HttpClient 更灵活,可适用于所有 Java 项目,它更底层,需要程序员自己管理和回收资源。在选择使用哪种工具时,可以根据实际项目情况来进行取舍。如果你已经使用了Spring框架,你可以考虑使用RestTemplate ,它将提供更好的集成和默认配置。如果你需要更灵活、底层控制程度更高的HTTP请求库, 你可能更愿意选择 HttpClient 。
-
大家好,这个月给大家带来的是5月份codeArts板块技术干货合集贴分享。希望可以给大家带来帮助,快捷便利的获取到自己想要的技术干货。 1.centos下安装redis6.0.2方法及踩坑记录。 https://bbs.huaweicloud.com/forum/thread-0236120881445478099-1-1.html 2.64位MFC调用32位DLL https://bbs.huaweicloud.com/forum/thread-0228120819107276100-1-1.html 3.用js+html,全网最火build出爱心【转】 https://bbs.huaweicloud.com/forum/thread-0236120798795045092-1-1.html 4.Django权限系统auth模块用法解读【转】 https://bbs.huaweicloud.com/forum/thread-0248120796948452094-1-1.html 5.Django日志logging的配置和自定义添加方式【转】 https://bbs.huaweicloud.com/forum/thread-0248120796668399093-1-1.html 6.Django如何利用uwsgi和nginx修改代码自动重启【转】 https://bbs.huaweicloud.com/forum/thread-0228120794336186093-1-1.html 7.Python实现监控一个程序的运行情况【转】 https://bbs.huaweicloud.com/forum/thread-02115120793876180082-1-1.html 8.使用Python和Scrapy实现抓取网站数据【转】 https://bbs.huaweicloud.com/forum/thread-0243120793550271100-1-1.html 9.浅谈一下关于Python对XML的解析【转】 https://bbs.huaweicloud.com/forum/thread-0260120793358226074-1-1.html 10.如何用python多次调用exe文件运行不同的结果【转】 https://bbs.huaweicloud.com/forum/thread-0264120793008226093-1-1.html 11.django和vue互传图片并进行处理和展示【转】 https://bbs.huaweicloud.com/forum/thread-0228120792806281092-1-1.html 12.django python 获取当天日期的方法【转】 https://bbs.huaweicloud.com/forum/thread-0243120792719244099-1-1.html 13.python之关于数组和列表的区别及说明【转】 https://bbs.huaweicloud.com/forum/thread-0243120792501298098-1-1.html 14.python数组和矩阵的用法解读【转】 https://bbs.huaweicloud.com/forum/thread-0260120792197826073-1-1.html 15.Python利用wxPython实现长文本处理【转】 https://bbs.huaweicloud.com/forum/thread-02115120792101273081-1-1.html 16.django如何计算两个TimeField的时差【转】 https://bbs.huaweicloud.com/forum/thread-0236120792022558090-1-1.html 17.python连接kafka加载数据的项目实践【转】 https://bbs.huaweicloud.com/forum/thread-0243120791900417097-1-1.html 18.数据翻译——Easy_Trans的简单使用 https://bbs.huaweicloud.com/forum/thread-0236120706211297083-1-1.html 19.vue如何使用print-nb库来实现打印功能 https://bbs.huaweicloud.com/forum/thread-02115120450806805064-1-1.html 20.如何使用Vue.js来实现打印功能 https://bbs.huaweicloud.com/forum/thread-0264120450582740073-1-1.html 21.使用Java Cron表达式实现定时器 https://bbs.huaweicloud.com/forum/thread-0243120448768906076-1-1.html 22.JDK8 Stream 操作 https://bbs.huaweicloud.com/forum/thread-0248120389213456072-1-1.html 23.Springboot在有锁的情况下如何正确使用事务 https://bbs.huaweicloud.com/forum/thread-0228120388866643073-1-1.html 24.免费漂亮的Java图形验证码 https://bbs.huaweicloud.com/forum/thread-02115120388590161062-1-1.html 25.Vue-CLI 3.0中配置production gzip https://bbs.huaweicloud.com/forum/thread-02115120388359055061-1-1.html
上滑加载中
推荐直播
-
华为云码道Agent集成与鸿蒙实战2026/08/11 周二 19:00-21:00
王一男-华为云码道产品规划专家;李炎-华为云码道产品专家;彭江敏-华为云鸿蒙端云一体化开发专家
本次直播带你解读华为云码道7月份产品新特性、新功能。更有专家演示码道Agent Space × 钉钉机器集成实战,从0到1打通消息通道;码道鸿蒙端云一体化实战,快速搭建员工签到系统。
回顾中 -
华为云开发者AI素养直播课·第五期2026/09/04 周五 16:00-18:00
林华鼎-华为云AI开发者运营负责人;蒋春阳-华为云AI开发者案例开发专家
本期直播内容: AI工具体验营 · 第5-8课连讲。Agent-Team 多智能体协作完成毕业设计实践
回顾中 -
华为云开发者AI素养ClassRoom·第六期2026/09/08 周二 19:00-20:00
樊渊-2026华为软件挑战赛冠军
高手来了:看软挑高手解析二维排样问题—从工业难题到算法突破
回顾中
热门标签