程序员必备的十大技能(进阶版)之网络与高并发原理(四)

简介: 教程来源 https://bncne.cn/ 本节详解高并发四大核心设计模式:①生产者-消费者(基于BlockingQueue与高性能Disruptor);②请求合并(批量处理降低IO压力);③背压控制(RxJava与信号量/滑动窗口限流);④线程封闭(ThreadLocal安全复用资源)。兼顾性能、可靠性与可维护性。

五、高并发系统设计模式

5.1 生产者-消费者模式

public class ProducerConsumerPattern {

    // 使用BlockingQueue实现
    public static class BlockingQueueDemo {
        private final BlockingQueue<Order> queue = new LinkedBlockingQueue<>(10000);

        // 生产者
        class Producer implements Runnable {
            @Override
            public void run() {
                while (true) {
                    Order order = createOrder();
                    // 阻塞直到有空位
                    queue.put(order);
                    // 非阻塞:queue.offer(order, 100, TimeUnit.MILLISECONDS);
                }
            }
        }

        // 消费者
        class Consumer implements Runnable {
            @Override
            public void run() {
                while (true) {
                    Order order = queue.take();  // 阻塞直到有数据
                    processOrder(order);
                }
            }
        }

        // 批量消费(提高吞吐量)
        class BatchConsumer implements Runnable {
            private final List<Order> buffer = new ArrayList<>();
            private final int batchSize = 100;

            @Override
            public void run() {
                while (true) {
                    // 使用drainTo批量获取
                    queue.drainTo(buffer, batchSize);
                    if (!buffer.isEmpty()) {
                        processBatch(buffer);
                        buffer.clear();
                    } else {
                        Thread.sleep(10);
                    }
                }
            }
        }
    }

    // 使用Disruptor(无锁环形缓冲区,高性能)
    public static class DisruptorDemo {
        // Disruptor基于RingBuffer,预分配内存,避免GC
        // 单个生产者可达千万级TPS

        // 定义事件
        static class OrderEvent {
            private long orderId;
            private long userId;
            private BigDecimal amount;
            // 对象复用(避免GC)

            void set(long orderId, long userId, BigDecimal amount) {
                this.orderId = orderId;
                this.userId = userId;
                this.amount = amount;
            }
        }

        // 事件工厂(预分配)
        static class OrderEventFactory implements EventFactory<OrderEvent> {
            @Override
            public OrderEvent newInstance() {
                return new OrderEvent();
            }
        }

        // 事件处理器
        static class OrderEventHandler implements EventHandler<OrderEvent> {
            @Override
            public void onEvent(OrderEvent event, long sequence, boolean endOfBatch) {
                // 处理订单
                processOrder(event);
            }
        }

        public void start() {
            // RingBuffer大小(必须是2的幂)
            int bufferSize = 1024 * 1024;

            Disruptor<OrderEvent> disruptor = new Disruptor<>(
                new OrderEventFactory(),
                bufferSize,
                DaemonThreadFactory.INSTANCE,
                ProducerType.MULTI,      // 多生产者
                new BusySpinWaitStrategy() // 忙等待策略(低延迟)
            );

            disruptor.handleEventsWith(new OrderEventHandler());
            disruptor.start();

            RingBuffer<OrderEvent> ringBuffer = disruptor.getRingBuffer();

            // 发布事件
            long sequence = ringBuffer.next();
            try {
                OrderEvent event = ringBuffer.get(sequence);
                event.set(123L, 456L, new BigDecimal("99.99"));
            } finally {
                ringBuffer.publish(sequence);
            }
        }
    }
}

5.2 请求合并(Request Coalescing)

@Component
public class RequestCoalescingService {

    // 将多个相同请求合并为一个批处理请求
    private final BlockingQueue<RequestPromise> queue = new LinkedBlockingQueue<>();
    private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();

    static class RequestPromise {
        final Long userId;
        final CompletableFuture<UserInfo> future;

        RequestPromise(Long userId, CompletableFuture<UserInfo> future) {
            this.userId = userId;
            this.future = future;
        }
    }

    @PostConstruct
    public void init() {
        // 定时批量处理(每10ms或积累到100个)
        scheduler.scheduleAtFixedRate(this::processBatch, 10, 10, TimeUnit.MILLISECONDS);
    }

    // 用户调用接口
    public CompletableFuture<UserInfo> getUserInfo(Long userId) {
        CompletableFuture<UserInfo> future = new CompletableFuture<>();
        queue.offer(new RequestPromise(userId, future));
        return future;
    }

    private void processBatch() {
        List<RequestPromise> batch = new ArrayList<>();
        queue.drainTo(batch, 100);  // 最多100个

        if (batch.isEmpty()) return;

        // 提取所有userId
        Set<Long> userIds = batch.stream()
            .map(rp -> rp.userId)
            .collect(Collectors.toSet());

        // 批量查询数据库(一次查询获取多个用户)
        Map<Long, UserInfo> userMap = userService.batchGetUsers(userIds);

        // 返回结果
        for (RequestPromise rp : batch) {
            UserInfo user = userMap.get(rp.userId);
            if (user != null) {
                rp.future.complete(user);
            } else {
                rp.future.completeExceptionally(new UserNotFoundException());
            }
        }
    }
}

5.3 背压(Backpressure)处理

public class BackpressureHandling {

    // 使用RxJava实现背压
    public void rxJavaBackpressure() {
        Flowable.range(1, 1_000_000)
            .onBackpressureBuffer(10000)           // 缓冲10000个
            // .onBackpressureDrop()               // 丢弃
            // .onBackpressureLatest()              // 只保留最新
            .observeOn(Schedulers.computation())
            .subscribe(
                value -> processSlowly(value),
                error -> log.error("Error", error),
                () -> log.info("Complete")
            );
    }

    // 自定义背压实现(速率限制)
    public static class RateLimitingProcessor {
        private final Semaphore semaphore;  // 信号量控制并发

        public RateLimitingProcessor(int maxConcurrent) {
            this.semaphore = new Semaphore(maxConcurrent);
        }

        public <T> CompletableFuture<T> process(Supplier<T> task) {
            CompletableFuture<T> future = new CompletableFuture<>();

            // 异步处理,等待许可
            CompletableFuture.runAsync(() -> {
                try {
                    semaphore.acquire();  // 背压:获取不到许可时阻塞
                    try {
                        T result = task.get();
                        future.complete(result);
                    } finally {
                        semaphore.release();
                    }
                } catch (InterruptedException e) {
                    future.completeExceptionally(e);
                    Thread.currentThread().interrupt();
                } catch (Exception e) {
                    future.completeExceptionally(e);
                }
            });

            return future;
        }

        // 使用滑动窗口控制速率
        public static class SlidingWindowRateLimiter {
            private final int maxRequests;
            private final long windowMillis;
            private final Queue<Long> timestamps = new ConcurrentLinkedQueue<>();

            public SlidingWindowRateLimiter(int maxRequests, long windowMillis) {
                this.maxRequests = maxRequests;
                this.windowMillis = windowMillis;
            }

            public synchronized boolean tryAcquire() {
                long now = System.currentTimeMillis();
                // 清理过期的请求
                while (!timestamps.isEmpty() && now - timestamps.peek() > windowMillis) {
                    timestamps.poll();
                }

                if (timestamps.size() < maxRequests) {
                    timestamps.offer(now);
                    return true;
                }
                return false;
            }
        }
    }
}

5.4 线程封闭与ThreadLocal

public class ThreadLocalUsage {

    // 1. 数据库连接管理
    public class ConnectionManager {
        private static final ThreadLocal<Connection> connectionHolder = new ThreadLocal<>() {
            @Override
            protected Connection initialValue() {
                return createConnection();
            }
        };

        public static Connection getConnection() {
            return connectionHolder.get();
        }

        public static void removeConnection() {
            connectionHolder.remove();  // 防止内存泄漏(线程池场景)
        }
    }

    // 2. 用户上下文传递
    public class UserContext {
        private static final ThreadLocal<User> currentUser = new ThreadLocal<>();

        public static void setCurrentUser(User user) {
            currentUser.set(user);
        }

        public static User getCurrentUser() {
            return currentUser.get();
        }

        // 子线程继承(InheritableThreadLocal)
        private static final InheritableThreadLocal<RequestId> requestId = new InheritableThreadLocal<>();

        // 线程池传递(使用阿里TransmittableThreadLocal)
        // TransmittableThreadLocal<String> context = new TransmittableThreadLocal<>();
        // TtlRunnable.get(originalRunnable) 包装任务
    }

    // 3. SimpleDateFormat线程安全问题(ThreadLocal包装)
    public class DateUtil {
        private static final ThreadLocal<SimpleDateFormat> dateFormatHolder = ThreadLocal.withInitial(
            () -> new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
        );

        public static String format(Date date) {
            return dateFormatHolder.get().format(date);
        }
    }

    // 4. 性能计数器(每个线程独立计数)
    public class PerThreadCounter {
        private static final ThreadLocal<LongAdder> counter = ThreadLocal.withInitial(LongAdder::new);

        public static void increment() {
            counter.get().increment();
        }

        public static long getAndReset() {
            long value = counter.get().sum();
            counter.get().reset();
            return value;
        }
    }
}

来源:
https://yvyus.cn/

相关文章
|
3月前
|
Java 程序员 测试技术
程序员必备的十大技能(进阶版)之网络与高并发原理(五)
教程来源 https://tmywi.cn/ 本节系统讲解高并发场景下的操作系统与JVM深度调优:涵盖Linux网络内核参数(BBR拥塞控制、缓冲区、连接队列)、进程资源限制、JVM G1GC优化及性能压测与瓶颈定位方法,并以百万级WebSocket推送系统为实战案例,体现无锁化、异步非阻塞等高性能设计原则。
|
Web App开发 编解码 前端开发
VUE网页实时播放海康、大华摄像头RTSP视频流完全方案,300毫秒延迟,支持H.265、可多路同时播放
在遍地都是摄像头的今天,往往需要在各种信息化、数字化、可视化B/S系统中集成实时视频流播放等功能,海康、大华、华为等厂家摄像头或录像机等设备一般也都遵循监控行业标准,支持国际标准的主流传输协议RTSP输出,而Chrome、Firefox、Edge等新一代浏览器从2015年开始取消了NPAPI插件技术支持导致RTSP流无法直接原生播放了
4547 0
|
3月前
|
人工智能 前端开发 Java
让 AI + Vercel 帮我部署网站,太方便了!
大家好,我是程序员鱼皮,今天分享一个提升 AI 编程效率的神技巧。 以前,我们开发一个网站,从写代码到部署上线,每个环节都得自己动手。 但是现在 AI 写代码已经溜飞边子了,在我心里已经取代了古法编程,谁能想到这是短短 1 年就发生的巨变。 而且 AI 的能力不止于此,利用 Skills 技能或 MCP 扩展,AI 甚至可以直接帮我们把网站部署上线! 好好好,合着我 只要提个一句话需求,写代码和部
225 0
|
8月前
|
敏捷开发 监控 数据可视化
产品研发轻量化管理工具(Sprint Board):敏捷落地的核心载体,让迭代效率倍增
Sprint Board 是面向敏捷团队的轻量化迭代管理工具,以极简看板串联“需求规划-任务拆解-执行跟踪-交付复盘”全流程。支持拖拽操作、实时同步、燃尽图与阻塞标记,助力中小团队快速落地Scrum,聚焦价值交付,降低协作内耗。(239字)
|
4月前
|
人工智能 Shell API
Claude Code 企业落地观察:近两天更新暴露的 MCP、代理、权限和模型网关问题
Claude Code 在 2026 年 5 月 8 日至 5 月 9 日连续更新,修复了 MCP OAuth、VS Code、Plan mode、代理链路和 Windows/WSL 体验问题。对企业团队来说,重点不是安装,而是治理。
549 1
|
5月前
|
人工智能 自然语言处理 安全
【适合新手的】5 分钟完成 OpenClaw Windows 一键部署实操
OpenClaw(小龙虾)是本地运行的AI智能体工具,支持自然语言驱动的文件整理、办公自动化与数据处理。本教程提供Windows一键部署方案:解压即用、无需代码、5分钟搞定,全程可视化操作,小白友好,安全高效。
|
5月前
|
传感器 人工智能 算法
智能猫砂盆如何实现“真正自动”?从智能算法看其技术路径
自动猫砂盆是否靠谱,关键在算法而非结构。
|
9月前
|
弹性计算 关系型数据库 数据库
阿里云卡券解析:优惠券、代金券、提货券、储值卡领取和使用指南及常见问题
为了助力更多新用户和老用户优惠上云,阿里云推出了多种优惠券、代金券、提货券和储值卡等多种卡券福利。这些券种不仅为用户提供了实实在在的优惠,还增加了购买阿里云产品的灵活性和便利性。本文将详细解析阿里云优惠券、代金券、提货券和储值卡的定义、用途、领取方式、使用规则及常见问题解答,以供大家了解他们之间的区别。
|
人工智能 云计算 UED
「云工开物」官网全面焕新,持续助力高校AI人才培养与科研创新
阿里云“云工开物”计划致力于让计算普惠高校师生,推动AI时代的人才培养与科研创新。过去一年多,60万大学生受益于AI和算力支持,300余所高校的人工智能课程得到助力,200余所高校开展AI实训。近期,“云工开物”官网焕新升级,新增七大功能板块,优化用户体验,提供一站式资源获取、学习中心、活动中心、教学与科研合作等服务,助力高校师生掌握AI技能、参与实践并加速科研创新。
|
Java 中间件 调度
SpringBoot整合XXL-JOB【03】- 执行器的使用
本文介绍了如何将调度中心与项目结合,通过配置“执行器”实现定时任务控制。首先新建SpringBoot项目并引入依赖,接着配置xxl-job相关参数,如调度中心地址、执行器名称等。然后通过Java代码将执行器注册为Spring Bean,并声明测试方法使用`@XxlJob`注解。最后,在调度中心配置并启动定时任务,验证任务是否按预期执行。通过这些步骤,读者可以掌握Xxl-Job的基本使用,专注于业务逻辑的编写而无需关心定时器本身的实现。
5514 10
SpringBoot整合XXL-JOB【03】-  执行器的使用