寻求适配公交站场景的Java并发同步代码(含问题代码)
Fixing Java Concurrency for Bus-Rider Synchronization Scenario
Let's walk through fixing your bus-rider concurrency implementation. First, let's restate the core requirements to make sure we're aligned, then break down what's wrong with your current code, and finally share a corrected version that meets all the scenario rules.
Core Scenario & Requirements
- When a bus arrives, all currently waiting riders board immediately; any riders arriving while the bus is loading must wait for the next bus.
- Buses have a capacity of 50 riders. If more than 50 riders are waiting, excess riders hold for the next bus.
- Buses depart right after all eligible riders board, or immediately if no riders are waiting.
- Buses arrive with exponential inter-arrival times (mean 20 minutes), riders arrive with exponential inter-arrival times (mean 30 seconds), and both run continuously.
What's Wrong With Your Existing Code
Your current code has several critical concurrency and logic gaps:
- One-Time Task Execution: Your
execute()method only submits a single rider and single bus task. We need continuous generation of both with the specified exponential delays. - Missing Synchronization Coordination: There's no way for riders to wait for a bus, or for buses to signal they're ready to load. The
isArrivedflag is never reset after a bus departs, so new riders will incorrectly be sent to the waiting stage forever after the first bus arrives. - Unsafe Collections:
ArrayListisn't thread-safe, so even with semaphores, adding/removing riders can cause race conditions or inconsistent state. - Abandoned Waiting Riders: Riders added to
waiting_stage_queueare never moved back to the bus stand when space opens up. - No Proper Bus Loading Logic: The code loads all riders in the bus stand, but doesn't enforce the 50-capacity rule properly (though your
putRiderlimits to 50, the loading step doesn't handle edge cases like partial loads).
Fixed Implementation
Let's rewrite the code with proper concurrency controls, continuous task scheduling, and correct state management.
1. Bus Class (With ID for Better Logging)
public class Bus { private final int id; public Bus(int id) { this.id = id; } public void depart() { System.out.printf("Bus %d is departing from the bus stand....%n", id); } public int getId() { return id; } }
2. Rider Class (With ID for Traceability)
public class Rider { private final int id; public Rider(int id) { this.id = id; } public void boardBus(Bus bus) { System.out.printf("Rider %d is boarding Bus %d%n", id, bus.getId()); } public int getId() { return id; } }
3. BusStandManager (Core Concurrency Logic)
We'll use thread-safe queues, semaphores for mutual exclusion, and a scheduled executor to handle continuous arrivals with exponential delays:
import java.util.concurrent.*; import java.util.Random; public class BusStandManager { private static final int BUS_CAPACITY = 50; // Bounded queue for riders waiting at the bus stand (max 50) private final BlockingQueue<Rider> busStandQueue = new LinkedBlockingQueue<>(BUS_CAPACITY); // Unbounded queue for riders who can't fit into the bus stand private final BlockingQueue<Rider> waitingStageQueue = new LinkedBlockingQueue<>(); // Ensures only one bus loads riders at a time private final Semaphore busMutex = new Semaphore(1); private final Random random = new Random(); private int nextBusId = 1; private int nextRiderId = 1; // Generate exponential inter-arrival time (converts mean seconds to ms) private long getExponentialDelay(double meanSeconds) { double meanMs = meanSeconds * 1000; return (long) (-meanMs * Math.log(1 - random.nextDouble())); } // Handles rider arrival: adds to bus stand if space exists, else waiting stage private void riderArrival() { Rider rider = new Rider(nextRiderId++); System.out.printf("Rider %d arrives at the bus stand%n", rider.getId()); try { // Try to add to bus stand (times out after 10ms to avoid hanging) if (!busStandQueue.offer(rider, 10, TimeUnit.MILLISECONDS)) { waitingStageQueue.put(rider); System.out.printf("Rider %d moves to waiting stage (bus stand full)%n", rider.getId()); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.printf("Rider %d thread interrupted%n", rider.getId()); } } // Handles bus arrival: loads riders, departs, then moves waiting riders to bus stand private void busArrival() { Bus bus = new Bus(nextBusId++); System.out.printf("Bus %d arrives at the bus stand%n", bus.getId()); try { busMutex.acquire(); // Lock to prevent concurrent bus loading int ridersToLoad = Math.min(busStandQueue.size(), BUS_CAPACITY); if (ridersToLoad > 0) { System.out.printf("Bus %d starts loading %d riders%n", bus.getId(), ridersToLoad); for (int i = 0; i < ridersToLoad; i++) { Rider rider = busStandQueue.take(); rider.boardBus(bus); } } else { System.out.printf("Bus %d finds no riders, preparing to depart%n", bus.getId()); } // Bus departs bus.depart(); // Fill bus stand with riders from waiting stage (up to capacity) int availableSlots = BUS_CAPACITY - busStandQueue.size(); if (availableSlots > 0 && !waitingStageQueue.isEmpty()) { int ridersToMove = Math.min(availableSlots, waitingStageQueue.size()); System.out.printf("Moving %d riders from waiting stage to bus stand%n", ridersToMove); for (int i = 0; i < ridersToMove; i++) { Rider rider = waitingStageQueue.take(); busStandQueue.put(rider); System.out.printf("Rider %d moves from waiting stage to bus stand%n", rider.getId()); } } busMutex.release(); // Unlock for next bus } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.printf("Bus %d thread interrupted%n", bus.getId()); } } // Starts the continuous simulation public void startSimulation() { ScheduledExecutorService executor = Executors.newScheduledThreadPool(4); // Schedule continuous rider arrivals (mean 30 seconds) executor.schedule(() -> { riderArrival(); // Reschedule next rider with exponential delay executor.schedule(this::riderArrival, getExponentialDelay(30), TimeUnit.MILLISECONDS); }, 0, TimeUnit.MILLISECONDS); // Schedule continuous bus arrivals (mean 20 minutes = 1200 seconds) executor.schedule(() -> { busArrival(); // Reschedule next bus with exponential delay executor.schedule(this::busArrival, getExponentialDelay(1200), TimeUnit.MILLISECONDS); }, 0, TimeUnit.MILLISECONDS); // Add shutdown hook to clean up executor on exit Runtime.getRuntime().addShutdownHook(new Thread(() -> { executor.shutdownNow(); System.out.println("Simulation stopped."); })); } }
4. Main Class
public class Main { public static void main(String[] args) { BusStandManager manager = new BusStandManager(); System.out.println("Starting bus-rider simulation..."); manager.startSimulation(); } }
Key Fixes Explained
- Thread-Safe Queues:
LinkedBlockingQueuehandles all thread-safe operations (add/remove) out of the box, eliminating race conditions from usingArrayList. The bus stand queue is bounded to 50, so riders are automatically routed to the waiting stage when it's full. - Continuous Arrivals: We use a scheduled executor to generate riders and buses with exponential inter-arrival times, matching the scenario's requirements for continuous operation.
- Mutual Exclusion: The
busMutexensures only one bus can load riders at a time, preventing conflicts where multiple buses try to load riders simultaneously. - Waiting Stage Handling: After a bus departs, we move riders from the waiting stage to fill the bus stand to capacity, ensuring no riders are stuck indefinitely.
- Proper Capacity Enforcement: Buses load up to their 50-rider capacity, and any excess riders wait for the next bus.
- Clean Shutdown: A shutdown hook ensures the executor is properly stopped when the program exits, preventing hanging threads.
内容的提问来源于stack exchange,提问作者samson j
相关产品推荐
相关产品推荐

