首页 文章 精选 留言 我的

精选列表

搜索[智能推荐],共10000篇文章
优秀的个人博客,低调大师

推荐!DevOps工具正越来越自动化

2009年,比利时根特市举办了首届 DevOpsDays 大会。至此,Development (开发)与 Operation (运维)的概念合二为一,被缩写为 DevOps (开发运维一体化)。 这一概念的风行并不在意料之外。亚马逊早期就提倡SOA (Service Oriented Architecture ),在亚马逊的每一个工程师都可以完全独立地完成编写代码,测试代码,版本管理,部署上线,服务监测等工作。现在,亚马逊凭借对 DevOps 文化的最佳实践,一跃成为世界级别IT领导者。 目前,DevOps 仍处于高速发展阶段,但如果做到整个业务部署 DevOps,不仅对软性的文化有要求,也有对硬性工具链的要求。南京大学软件研发效能实验室发布的《DevOps ·云原生2021年度中国调查报告》显示,2021年国内企业的 DevOps 工具的普及程度较2019年有明显上升。 DevOps 工具的使用正在变得越来越重要。随着各色工具的出现与应用,让 DevOps 得以从一个概念慢慢变成现实。 目前,DevOps 工具覆盖了从规划、编码、构建、测试、发布、部署和维护的软件生产全过程。 01 敏捷开发工具,加速开发效率 “敏捷”(Agile)这个概念起源于美国,比 DevOps 出现的时间要早。 敏捷开发框架被提出以后,软件的运行和运维方式发生了巨大的变化。因为“敏捷”完全摒弃了传统的瀑布式思维,将小步快跑、不断迭代的分步思想融汇在开发过程中。 随后,敏捷成为 DevOps 中的重要一环,一些敏捷教练认为,工作在这个系统中的所有人需要有全局协作和优化的意识,优化价值流的流量、流速,关注价值的产生。工程师们更要有工具链整体优化的意识,而不仅仅是精通某个环节,或局限在与自己相关的上下游工具上。 在敏捷开发中,往往离不开一些需求和项目管理工具,比如 JIRA、Trello等。 JIRA是一个被广泛使用的问题跟踪器,提供了 bug 跟踪和敏捷项目管理功能。虽然JIRA是由澳大利亚 Atlassian 公司创建的商业授权产品,但也有可有限使用的免费版本。 JIRA详情可查看:https://www.oschina.net/p/jira 此外,Atlassian 子公司所拥有的 Trello 也是在国际上非常流行的协作工具,以设计简约著名,很多团队用它来计划各自的工作 sprint(冲刺)。国内知名的Teambition、Leangoo等都在一定程度上效仿了Trello的设计。 02 持续交付工具,助力落地DevOps DevOps 是一个完整的面向IT运维的工作流,其中CI/CD是基础。 CI是Continuous Integration(持续集成),而CD对应多个英文,Continuous Delivery(持续交付)或Continuous Deployment(持续部署)。 如果无法做到CI/CD的话,DevOps 也就变成了空中楼阁。 要做到持续交付,构建十分重要。在构建阶段,需要保持打包的一致性,而自动执行容易出错的活动,生成早期质量信号。此时,Maven、NPM、Gradle等构建工具则大有用处。 Maven 主要用于 Java 项目的自动化构建,同时它也可以构建和管理以 C#、Ruby、Scala 等语言编写的项目。而Gradle 则是基于 Apache Ant 和 Apache Maven 理念的自动化构建系统,它还引入了基于 Groovy 的领域特定语言。 Maven详情可查看:https://www.oschina.net/p/maven Gradle详情可查看:https://www.oschina.net/p/gradle 在构建Maven库前后,往往需要Nexus这样的仓库管理工具。Nexus 是一套“开箱即用”的系统,功能非常强大,它极大简化了自己内部仓库的维护和外部仓库的访问。在内部,你可以配置构建工具,并发布到 Nexus,然后其他开发人员就可以使用它们了。 Nexus 详情可查看:https://www.oschina.net/p/nexus 在持续集成中,像Jenkins、Bamboo等这样的流水线工具不可或缺。其中,Jenkins 是开源、免费、与平台无关的自动化服务器,它独立于 Java ,且支持Windows、Mac 和其他类似 UNIX 的操作系统,通过Jenkins 可以将本机系统软件包 Docker 安装。 Jenkins 详情可查看:https://www.oschina.net/p/jenkins 此外,容器引擎和编排工具也在DevOps 的持续交付工具之列。南京大学软件研发效能实验室发布的《DevOps·云原生2021年度中国调查报告》显示,过去两年,容器技术的应用持续深化,以容器及其编排技术为核心的生态,逐渐扩展至涵盖微服务、DevOps、服务监测分析、应用管理的完整闭环。 因此,以Docker、K8s、Apache Mesos等为代表的容器引擎和编排工具几乎在DevOps实践中扮演着不可替代的角色。 Apache Mesos详情可查看:https://www.oschina.net/p/apache+mesos 03 自动化运维工具,补上运维一环 DevOps将传统的“开发”和“运维”概念合二为一,将两者从传统作坊式的工作方式解放出来。所以,在DevOps 工具链中,Zabbix、Elastic、Grafana、Kafka、Ansible、Logstash、Prometheus等自动运维工具的作用不可小觑。 其中,Ansible 是由 RedHat 维护的开源 IT 自动化工具,使用剧本(playbooks)做配置管理和多机部署系统,它运行在 Unix 家族系统上,可以配置 Unix 家族系统和 Windows。 我们可以在控制机器上安装 Ansible,而不需要 Ansible 在其他服务器上运行,这些服务器可以从 Web 到应用程序再到数据库服务器。 Ansible 详情可查看:https://www.oschina.net/p/ansible Prometheus则是用于事件监视和警报的免费软件应用程序,它在时间序列数据库中记录实时指标,基于 HTTP 拉取模型,支持灵活的查询和实时警报。Prometheus 服务器的工作方式是抓取,也就是说,调用各个节点暴露出来的指标端点。 Prometheus 详情可查看:https://www.oschina.net/p/prometheus Grafana包括企业版和开源版本两种,是可视化分析软件,可以查询、可视化、报警和探索指标,无论这些指标存储在哪里,Grafana 都可以通过提供相关数据来帮助我们跟踪用户行为、应用程序行为、在生产环境或预生产环境中弹出错误的频率、弹出错误的类型以及上下文场景。 Grafana 详情可查看:https://www.oschina.net/p/grafana 04 平台类工具,DevOps工具的“集大成者” 随着 DevOps 实践在国内外企业中流行开来,用户对自动化的要求越来越高。因此,也就催生了更多集成功能的DevOps 平台,例如我国的“飞算 SoFlu全自动软件工程平台”。 飞算 SoFlu其实是通过可视化编程的方式满足开发需求,也就是说,输入流程图即可实现自动开发、自动测试、自动运维等,由此提高工作效率,使用户可以更多关注自身业务。在平台使用过程中,可以达到一个ID相当于一个10人科技团队的效果。 此外,飞算还可以通过管理平台来管理需求、研发、测试、部署、上线、运维等整个软件生命周期,沉淀经验、积累知识,将管理制度真正落地。 以近期上线的测试平台为例,飞算SoFlu 通过自动化的生命周期管理、测试用例自动生成、测试数据管理等功能,去解决人工测试耗时长、测试跟踪管理难、测试成本高等难题。软件质量可以通过工具、流程和管理予以保障,而不再依靠有丰富经验的软件工程师。 目前,飞算 SoFlu仍在加速更新更加强大的 DevOps 功能,未来可期。 飞算详情可查看:https://www.feisuanyz.com/ 05 微服务相关技术,通过拆分更加便利 所谓“微服务”,就是将原来黑盒化的一个整体产品进行拆分(解耦),从一个提供多种服务的整体,拆成各自提供不同服务的多个个体。因此,通过微服务技术,不同的工程师可以对各自负责的模块进行处理,例如开发、测试、部署、迭代。 Spring cloud、Spring Boot、Apache Dubbo等工具都可以用于微服务之中。其中,Spring cloud 专注为典型的用例和可扩展性机制提供良好的开箱即用体验,它为开发人员提供了快速构建分布式系统中常见模式的功能,通过Spring cloud 开发者可以快速实现样板的服务和应用程序。 Spring cloud 详情可查看:https://www.oschina.net/p/spring-cloud 06 安全管理工具,让DevOps 走向DevSecOps DevSecOps,也就是研发安全运营一体化,将安全融入 DevOps 每个阶段过程,开发、安全、运营各部门紧密合作,强调在安全风险可控的前提下,帮助企业提升IT效能,更好地实现研发运营一体化。 云计算开源产业联盟发布的《中国DevOps现状调查报告2021》显示,源代码静态安全检测、容器镜像安全扫描及 Web 应用防火墙(WAF)正在成为企业应用最广泛的 DevSecOps 技术实践。 并且,企业在选择 DevOps 工具时更注重功能的易用性、工具自身的安全性和自动化程度。 调查显示, 超过四成的企业在选择 DevOps 工具时考虑工具的功能的易用性(43.18%)、工具自身的安全性 (42.96%)和工具的自动化程度(42.80%)。 安全工具也是百花齐放, 包括了代码安全工具Fortify、容器安全工具Clair、Web安全工具AppScan等等。

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

博文推荐|Pulsar 客户端编码最佳实践

本文描述了一些 Pulsar 客户端编码相关的最佳实践,并提供了可商用的样例代码,供大家研发的时候参考,提升大家接入 Pulsar 的效率。在生产环境上,Pulsar 的地址信息往往都通过配置中心或者是 K8s 域名发现的方式获得,这块不是这篇文章描述的重点,以PulsarConstant.SERVICE_HTTP_URL代替。本文中的例子均已上传到Github[1]。 前期 Client 初始化和配置 初始化 Client--demo 级别 import lombok.extern.slf4j.Slf4j;import org.apache.pulsar.client.api.PulsarClient;/*** @author hezhangjian*/@Slf4jpublic class DemoPulsarClientInit { private static final DemoPulsarClientInit INSTANCE = new DemoPulsarClientInit(); private PulsarClient pulsarClient; public static DemoPulsarClientInit getInstance() { return INSTANCE; } public void init() throws Exception { pulsarClient = PulsarClient.builder() .serviceUrl(PulsarConstant.SERVICE_HTTP_URL) .build(); } public PulsarClient getPulsarClient() { return pulsarClient; }} Demo 级别的 Pulsar client 初始化的时候没有配置任何自定义参数,并且初始化的时候没有考虑异常,init的时候会直接抛出异常。 初始化 Client--可上线级别 import io.netty.util.concurrent.DefaultThreadFactory;import lombok.extern.slf4j.Slf4j;import org.apache.pulsar.client.api.PulsarClient;import java.util.concurrent.Executors;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.TimeUnit;/*** @author hezhangjian*/@Slf4jpublic class DemoPulsarClientInitRetry { private static final DemoPulsarClientInitRetry INSTANCE = new DemoPulsarClientInitRetry(); private volatile PulsarClient pulsarClient; private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-cli-init")); public static DemoPulsarClientInitRetry getInstance() { return INSTANCE; } public void init() { executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS); } private void initWithRetry() { try { pulsarClient = PulsarClient.builder() .serviceUrl(PulsarConstant.SERVICE_HTTP_URL) .build(); log.info("pulsar client init success"); this.executorService.shutdown(); } catch (Exception e) { log.error("init pulsar error, exception is ", e); } } public PulsarClient getPulsarClient() { return pulsarClient; }} 在实际的环境中,我们往往要做到pulsar client初始化失败后不影响微服务的启动,即待微服务启动后,再一直重试创建pulsar client。上面的代码示例通过volatile加不断循环重建实现了这一目标,并且在客户端成功创建后,销毁了定时器线程。 初始化 Client--商用级别 import io.netty.util.concurrent.DefaultThreadFactory;import lombok.extern.slf4j.Slf4j;import org.apache.pulsar.client.api.PulsarClient;import org.apache.pulsar.client.api.SizeUnit;import java.util.concurrent.Executors;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.TimeUnit;/*** @author hezhangjian*/@Slf4jpublic class DemoPulsarClientInitUltimate { private static final DemoPulsarClientInitUltimate INSTANCE = new DemoPulsarClientInitUltimate(); private volatile PulsarClient pulsarClient; private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-cli-init")); public static DemoPulsarClientInitUltimate getInstance() { return INSTANCE; } public void init() { executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS); } private void initWithRetry() { try { pulsarClient = PulsarClient.builder() .serviceUrl(PulsarConstant.SERVICE_HTTP_URL) .ioThreads(4) .listenerThreads(10) .memoryLimit(64, SizeUnit.MEGA_BYTES) .operationTimeout(5, TimeUnit.SECONDS) .connectionTimeout(15, TimeUnit.SECONDS) .build(); log.info("pulsar client init success"); this.executorService.shutdown(); } catch (Exception e) { log.error("init pulsar error, exception is ", e); } } public PulsarClient getPulsarClient() { return pulsarClient; }} 商用级别的Pulsar Client新增了 5 个配置参数: •ioThreadsnetty 的 ioThreads 负责网络 IO 操作,如果业务流量较大,可以调高ioThreads个数;•listenersThreads负责调用以listener模式启动的消费者的回调函数,建议配置大于该 client 负责的partition数目;•memoryLimit当前用于限制pulsar生产者可用的最大内存,可以很好地防止网络中断、Pulsar 故障等场景下,消息积压在producer侧,导致 Java 程序 OOM;•operationTimeout一些元数据操作的超时时间,Pulsar 默认为 30s,有些保守,可以根据自己的网络情况、处理性能来适当调低;•connectionTimeout连接 Pulsar 的超时时间,配置原则同上。 客户端进阶参数(内存分配相关) 我们还可以通过传递 Java 的 property 来控制 Pulsar 客户端内存分配的参数,这里列举几个重要参数: •pulsar.allocator.pooled为 true 则使用堆外内存池,false 则使用堆内存分配,不走内存池。默认使用高效的堆外内存池;•pulsar.allocator.exit_on_oom如果内存溢出,是否关闭jvm,默认为 false;•pulsar.allocator.out_of_memory_policy在 https://github.com/apache/pulsar/pull/12200 引入,目前还没有正式 release 版本,用于配置当堆外内存不够使用时的行为,可选项为FallbackToHeap和ThrowException,默认为FallbackToHeap,如果你不希望消息序列化的内存影响到堆内存分配,则可以配置成ThrowException。 生产者 初始化 producer 重要参数 maxPendingMessages 生产者消息发送队列,根据实际 topic 的量级合理配置,避免在网络中断、Pulsar 故障场景下的 OOM。建议和 client 侧的配置memoryLimit之间挑一个进行配置。 messageRoutingMode 消息路由模式。默认为RoundRobinPartition。根据业务需求选择,如果需要保序,一般选择SinglePartition,把相同 key 的消息发到同一个partition。 autoUpdatePartition 自动更新 partition 信息。如topic中partition信息不变则不需要配置,降低集群的消耗。 batch 相关参数 因为批量发送模式底层由定时任务实现,如果该 topic 上消息数较小,则不建议开启batch。尤其是大量的低时间间隔的定时任务会导致 netty 线程 CPU 飙高。 •enableBatching是否启用批量发送;•batchingMaxMessages批量发送最大消息条数;•batchingMaxPublishDelay批量发送定时任务间隔。 静态 producer 初始化 静态 producer,指不会随着业务的变化进行 producer 的启动或关闭。那么就在微服务启动完成、client 初始化完成之后,初始化 producer,样例如下: 一个生产者一个线程,适用于生产者数目较少的场景 import io.netty.util.concurrent.DefaultThreadFactory;import lombok.extern.slf4j.Slf4j;import org.apache.pulsar.client.api.Producer;import java.util.concurrent.Executors;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.TimeUnit;/*** @author hezhangjian*/@Slf4jpublic class DemoPulsarStaticProducerInit { private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-producer-init")); private final String topic; private volatile Producer<byte[]> producer; public DemoPulsarStaticProducerInit(String topic) { this.topic = topic; } public void init() { executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS); } private void initWithRetry() { try { final DemoPulsarClientInit instance = DemoPulsarClientInit.getInstance(); producer = instance.getPulsarClient().newProducer().topic(topic).create(); } catch (Exception e) { log.error("init pulsar producer error, exception is ", e); } } public Producer<byte[]> getProducer() { return producer; }} 多个生产者一个线程,适用于生产者数目较多的场景 import io.netty.util.concurrent.DefaultThreadFactory;import lombok.extern.slf4j.Slf4j;import org.apache.pulsar.client.api.Producer;import java.util.List;import java.util.concurrent.CopyOnWriteArrayList;import java.util.concurrent.Executors;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.TimeUnit;/*** @author hezhangjian*/@Slf4jpublic class DemoPulsarStaticProducersInit { private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-consumer-init")); private CopyOnWriteArrayList<Producer<byte[]>> producers; private int initIndex; private List<String> topics; public DemoPulsarStaticProducersInit(List<String> topics) { this.topics = topics; } public void init() { executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS); } private void initWithRetry() { if (initIndex == topics.size()) { return; } for (; initIndex < topics.size(); initIndex++) { try { final DemoPulsarClientInit instance = DemoPulsarClientInit.getInstance(); final Producer<byte[]> producer = instance.getPulsarClient().newProducer().topic(topics.get(initIndex)).create();; producers.add(producer); } catch (Exception e) { log.error("init pulsar producer error, exception is ", e); break; } } } public CopyOnWriteArrayList<Producer<byte[]>> getProducers() { return producers; }} 动态生成销毁的 producer 示例 还有一些业务,我们的 producer 可能会根据业务来进行动态的启动或销毁,如接收道路上车辆的数据并发送给指定的 topic。我们不会让内存里面驻留所有的 producer,这会导致占用大量的内存,我们可以采用类似于 LRU Cache 的方式来管理 producer 的生命周期。 /*** @author hezhangjian*/@Slf4jpublic class DemoPulsarDynamicProducerInit { /** * topic -- producer */ private AsyncLoadingCache<String, Producer<byte[]>> producerCache; public DemoPulsarDynamicProducerInit() { this.producerCache = Caffeine.newBuilder() .expireAfterAccess(600, TimeUnit.SECONDS) .maximumSize(3000) .removalListener((RemovalListener<String, Producer<byte[]>>) (topic, value, cause) -> { log.info("topic {} cache removed, because of {}", topic, cause); try { value.close(); } catch (Exception e) { log.error("close failed, ", e); } }) .buildAsync(new AsyncCacheLoader<>() { @Override public CompletableFuture<Producer<byte[]>> asyncLoad(String topic, Executor executor) { return acquireFuture(topic); } @Override public CompletableFuture<Producer<byte[]>> asyncReload(String topic, Producer<byte[]> oldValue, Executor executor) { return acquireFuture(topic); } }); } private CompletableFuture<Producer<byte[]>> acquireFuture(String topic) { CompletableFuture<Producer<byte[]>> future = new CompletableFuture<>(); try { ProducerBuilder<byte[]> builder = DemoPulsarClientInit.getInstance().getPulsarClient().newProducer().enableBatching(true); final Producer<byte[]> producer = builder.topic(topic).create(); future.complete(producer); } catch (Exception e) { log.error("create producer exception ", e); future.completeExceptionally(e); } return future; }} 这个模式下,可以根据返回的CompletableFuture<Producer<byte[]>>来优雅地进行流式处理。 可以接受消息丢失的发送 final CompletableFuture<Producer<byte[]>> cacheFuture = producerCache.get(topic); cacheFuture.whenComplete((producer, e) -> { if (e != null) { log.error("create pulsar client exception ", e); return; } try { producer.sendAsync(msg).whenComplete(((messageId, throwable) -> { if (throwable != null) { log.error("send producer msg error ", throwable); return; } log.info("topic {} send success, msg id is {}", topic, messageId); })); } catch (Exception ex) { log.error("send async failed ", ex); } }); 以上为正确处理Client创建失败和发送失败的回调函数。但是由于在生产环境下,Pulsar 并不是一直保持可用的,会因为虚拟机故障、Pulsar 服务升级等导致发送失败。这个时候如果要保证消息发送成功,就需要对消息发送进行重试。 可以容忍极端场景下的发送丢失 final Timer timer = new HashedWheelTimer(); private void sendMsgWithRetry(String topic, byte[] msg, int retryTimes) { final CompletableFuture<Producer<byte[]>> cacheFuture = producerCache.get(topic); cacheFuture.whenComplete((producer, e) -> { if (e != null) { log.error("create pulsar client exception ", e); return; } try { producer.sendAsync(msg).whenComplete(((messageId, throwable) -> { if (throwable == null) { log.info("topic {} send success, msg id is {}", topic, messageId); return; } if (retryTimes == 0) { timer.newTimeout(timeout -> DemoPulsarDynamicProducerInit.this.sendMsgWithRetry(topic, msg, retryTimes - 1), 1 << retryTimes, TimeUnit.SECONDS); } log.error("send producer msg error ", throwable); })); } catch (Exception ex) { log.error("send async failed ", ex); } }); } 这里在发送失败后,做了退避重试,可以容忍pulsar服务端故障一段时间。比如退避 7 次、初次间隔为 1s,那么就可以容忍1+2+4+8+16+32+64=127s的故障。这已经足够满足大部分生产环境的要求了。因为理论上存在超过 127s 的故障,所以还是要在极端场景下,向上游返回失败。 生产者 Partition 级别严格保序 生产者严格保序的要点:一次只发送一条消息,确认发送成功后再发送下一条消息。实现上可以使用同步异步两种模式: •同步模式的要点就是循环发送,直到上一条消息发送成功后,再启动下一条消息发送;•异步模式的要点是观测上一条消息发送的 future,如果失败也一直重试,成功则启动下一条消息发送。 值得一提的是,这个模式下,partition 间是可以并行的,可以使用OrderedExecutor或per partition per thread。 同步模式举例: /*** @author hezhangjian*/@Slf4jpublic class DemoPulsarProducerSyncStrictlyOrdered { Producer<byte[]> producer; public void sendMsg(byte[] msg) { while (true) { try { final MessageId messageId = producer.send(msg); log.info("topic {} send success, msg id is {}", producer.getTopic(), messageId); break; } catch (Exception e) { log.error("exception is ", e); } } }} 消费者 初始化消费者重要参数 receiverQueueSize 注意:处理不过来时,消费缓冲队列会积压在内存中,合理配置防止 OOM。 autoUpdatePartition 自动更新 partition 信息。如topic中partition信息不变则不需要配置,降低集群的消耗。 subscribeType 订阅类型,根据业务需求决定。 subscriptionInitialPosition 订阅开始的位置,根据业务需求决定最前或者最后。 messageListener 使用 listener 模式消费,只需要提供回调函数,不需要主动执行receive()拉取。一般没有特殊诉求,建议采用 listener 模式。 ackTimeout 当服务端推送消息,但消费者未及时回复 ack 时,经过 ackTimeout 后,会重新推送给消费者处理,即redeliver机制。注意在利用redeliver机制的时候,一定要注意仅仅使用重试机制来重试可恢复的错误。举个例子,如果代码里面对消息进行解码,解码失败就不适合利用redeliver机制。这会导致客户端一直处于重试之中。 如果拿捏不准,还可以通过下面的deadLetterPolicy配置死信队列,防止消息一直重试。 negativeAckRedeliveryDelay 当客户端调用negativeAcknowledge时,触发redeliver机制的时间。redeliver机制的注意点同ackTimeout。 需要注意的是,ackTimeout和negativeAckRedeliveryDelay建议不要同时使用,一般建议使用negativeAck,用户可以有更灵活的控制权。一旦ackTimeout配置的不合理,在消费时间不确定的情况下可能会导致消息不必要的重试。 deadLetterPolicy 配置redeliver的最大次数和死信 topic。 初始化消费者原则 消费者只有创建成功才能工作,不像生产者可以向上游返回失败,所以消费者要一直重试创建。示例代码如下:注意:消费者和 topic 可以是一对多的关系,消费者可以订阅多个 topic。 一个消费者一个线程,适用于消费者数目较少的场景 import io.netty.util.concurrent.DefaultThreadFactory;import lombok.extern.slf4j.Slf4j;import org.apache.pulsar.client.api.Consumer;import java.util.concurrent.Executors;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.TimeUnit;/*** @author hezhangjian*/@Slf4jpublic class DemoPulsarConsumerInit { private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-consumer-init")); private final String topic; private volatile Consumer<byte[]> consumer; public DemoPulsarConsumerInit(String topic) { this.topic = topic; } public void init() { executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS); } private void initWithRetry() { try { final DemoPulsarClientInit instance = DemoPulsarClientInit.getInstance(); consumer = instance.getPulsarClient().newConsumer().topic(topic).messageListener(new DemoMessageListener<>()).subscribe(); } catch (Exception e) { log.error("init pulsar producer error, exception is ", e); } } public Consumer<byte[]> getConsumer() { return consumer; }} 多个消费者一个线程,适用于消费者数目较多的场景 import io.netty.util.concurrent.DefaultThreadFactory;import lombok.extern.slf4j.Slf4j;import org.apache.pulsar.client.api.Consumer;import java.util.List;import java.util.concurrent.CopyOnWriteArrayList;import java.util.concurrent.Executors;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.TimeUnit;/*** @author hezhangjian*/@Slf4jpublic class DemoPulsarConsumersInit { private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new DefaultThreadFactory("pulsar-consumer-init")); private CopyOnWriteArrayList<Consumer<byte[]>> consumers; private int initIndex; private List<String> topics; public DemoPulsarConsumersInit(List<String> topics) { this.topics = topics; } public void init() { executorService.scheduleWithFixedDelay(this::initWithRetry, 0, 10, TimeUnit.SECONDS); } private void initWithRetry() { if (initIndex == topics.size()) { return; } for (; initIndex < topics.size(); initIndex++) { try { final DemoPulsarClientInit instance = DemoPulsarClientInit.getInstance(); final Consumer<byte[]> consumer = instance.getPulsarClient().newConsumer().topic(topics.get(initIndex)).messageListener(new DemoMessageListener<>()).subscribe(); consumers.add(consumer); } catch (Exception e) { log.error("init pulsar producer error, exception is ", e); break; } } } public CopyOnWriteArrayList<Consumer<byte[]>> getConsumers() { return consumers; }} 消费者达到至少一次语义 使用手动回复 ack 模式,确保处理成功后再 ack。如果处理失败可以自己重试或通过negativeAck机制进行重试 同步模式举例 这里需要注意,如果处理消息时长差距比较大,同步处理的方式可能会让本来可以很快处理的消息得不到处理的机会。 /*** @author hezhangjian*/@Slf4jpublic class DemoMessageListenerSyncAtLeastOnce<T> implements MessageListener<T> { @Override public void received(Consumer<T> consumer, Message<T> msg) { try { final boolean result = syncPayload(msg.getData()); if (result) { consumer.acknowledgeAsync(msg); } else { consumer.negativeAcknowledge(msg); } } catch (Exception e) { // 业务方法可能会抛出异常 log.error("exception is ", e); consumer.negativeAcknowledge(msg); } } /** * 模拟同步执行的业务方法 * @param msg 消息体内容 * @return */ private boolean syncPayload(byte[] msg) { return System.currentTimeMillis() % 2 == 0; }} 异步模式举例 异步的话需要考虑内存的限制,因为异步的方式可以很快地从broker消费,不会被业务操作阻塞,这样inflight的消息可能会非常多。如果是Shared或KeyShared模式,可以通过maxUnAckedMessage进行限制。如果是Failover模式,可以通过下面的消费者繁忙时阻塞拉取消息,不再进行业务处理通过判断inflight消息数来阻塞处理。 /*** @author hezhangjian*/@Slf4jpublic class DemoMessageListenerAsyncAtLeastOnce<T> implements MessageListener<T> { @Override public void received(Consumer<T> consumer, Message<T> msg) { try { asyncPayload(msg.getData(), new DemoSendCallback() { @Override public void callback(Exception e) { if (e == null) { consumer.acknowledgeAsync(msg); } else { log.error("exception is ", e); consumer.negativeAcknowledge(msg); } } }); } catch (Exception e) { // 业务方法可能会抛出异常 consumer.negativeAcknowledge(msg); } } /** * 模拟异步执行的业务方法 * @param msg 消息体 * @param demoSendCallback 异步函数的callback */ private void asyncPayload(byte[] msg, DemoSendCallback demoSendCallback) { if (System.currentTimeMillis() % 2 == 0) { demoSendCallback.callback(null); } else { demoSendCallback.callback(new Exception("exception")); } }} 消费者繁忙时阻塞拉取消息,不再进行业务处理 当消费者处理不过来时,通过阻塞listener方法,不再进行业务处理。避免在微服务积累太多消息导致 OOM,可以通过 RateLimiter 或者 Semaphore 控制处理。 *** @author hezhangjian*/@Slf4jpublic class DemoMessageListenerBlockListener<T> implements MessageListener<T> { /** * Semaphore保证最多同时处理500条消息 */ private final Semaphore semaphore = new Semaphore(500); @Override public void received(Consumer<T> consumer, Message<T> msg) { try { semaphore.acquire(); asyncPayload(msg.getData(), new DemoSendCallback() { @Override public void callback(Exception e) { semaphore.release(); if (e == null) { consumer.acknowledgeAsync(msg); } else { log.error("exception is ", e); consumer.negativeAcknowledge(msg); } } }); } catch (Exception e) { semaphore.release(); // 业务方法可能会抛出异常 consumer.negativeAcknowledge(msg); } } /** * 模拟异步执行的业务方法 * @param msg 消息体 * @param demoSendCallback 异步函数的callback */ private void asyncPayload(byte[] msg, DemoSendCallback demoSendCallback) { if (System.currentTimeMillis() % 2 == 0) { demoSendCallback.callback(null); } else { demoSendCallback.callback(new Exception("exception")); } }} 消费者严格按 partition 保序 为了实现partition级别消费者的严格保序,需要对单partition的消息,一旦处理失败,在这条消息重试成功之前不能处理该partition的其他消息。示例如下: /*** @author hezhangjian*/@Slf4jpublic class DemoMessageListenerSyncAtLeastOnceStrictlyOrdered<T> implements MessageListener<T> { @Override public void received(Consumer<T> consumer, Message<T> msg) { retryUntilSuccess(msg.getData()); consumer.acknowledgeAsync(msg); } private void retryUntilSuccess(byte[] msg) { while (true) { try { final boolean result = syncPayload(msg); if (result) { break; } } catch (Exception e) { log.error("exception is ", e); } } } /** * 模拟同步执行的业务方法 * * @param msg 消息体内容 * @return */ private boolean syncPayload(byte[] msg) { return System.currentTimeMillis() % 2 == 0; }} 致谢 感谢鹏辉[2]和罗天[3]的审稿。 作者简介 贺张俭[4],Apache Pulsar Contributor,西安电子科技大学毕业,华为云物联网高级工程师,目前 Pulsar 已经在华为云物联网大规模商用,了解更多内容可以访问他的简书博客地址[5]。 相关链接 •最佳实践|Apache Pulsar 在华为云物联网之旅 关注「Apache Pulsar」👇🏻,获取干货与动态 👇🏻 加入 Apache Pulsar 中文交流群 👇🏻 引用链接 [1]Github:https://github.com/Shoothzj/pulsar-client-examples[2]鹏辉:https://github.com/codelipenghui[3]罗天:https://github.com/fu-turer[4]贺张俭:https://github.com/Shoothzj[5]简书博客地址:https://www.jianshu.com/u/9e21abacd418 点击「阅读原文」,查看 Apache Pulsar 干货集锦 本文分享自微信公众号 - ApachePulsar(ApachePulsar)。如有侵权,请联系 support@oschina.cn 删除。本文参与“OSC源创计划”,欢迎正在阅读的你也加入,一起分享。

资源下载

更多资源
Spring

Spring

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

Rocky Linux

Rocky Linux

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

Sublime Text

Sublime Text

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

WebStorm

WebStorm

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

用户登录
用户注册