如何在指定时段停用Apache Storm拓扑并到期自动启动
Great question! Handling scheduled maintenance windows for your Apache Storm topology (that feeds data to an API) while integrating with your Spring Boot stack is totally feasible. Let’s break down practical, actionable approaches tailored to your tech stack:
This is the quickest approach if you’re already comfortable with Storm’s command-line tools. Here’s how to wire it up:
Step 1: Fetch maintenance windows from your service
Create a Spring component to call your maintenance window service regularly. Parse the response into start/end timestamps or cron expressions that represent when to stop/start the topology.@Service public class MaintenanceWindowService { private final RestTemplate restTemplate; public MaintenanceWindowService(RestTemplate restTemplate) { this.restTemplate = restTemplate; } public MaintenanceWindow getCurrentMaintenanceWindow() { // Call your existing maintenance window service return restTemplate.getForObject("/api/maintenance/window", MaintenanceWindow.class); } } // DTO for maintenance window record MaintenanceWindow(LocalDateTime start, LocalDateTime end, boolean isActive) {}Step 2: Schedule topology stop/start tasks
Use Spring’s@Scheduledannotation to run checks, or use Quartz if you need more flexible scheduling (like dynamic cron based on maintenance windows). For simple start/end triggers:@Service public class TopologySchedulerService { private final MaintenanceWindowService maintenanceService; private final TopologyCommandExecutor commandExecutor; public TopologySchedulerService(MaintenanceWindowService maintenanceService, TopologyCommandExecutor commandExecutor) { this.maintenanceService = maintenanceService; this.commandExecutor = commandExecutor; } // Check every minute if we need to act @Scheduled(fixedRate = 60000) public void checkMaintenanceStatus() { MaintenanceWindow window = maintenanceService.getCurrentMaintenanceWindow(); LocalDateTime now = LocalDateTime.now(); if (window.isActive() && now.isAfter(window.start()) && now.isBefore(window.end())) { // Stop the topology if it's running commandExecutor.stopTopology("your-topology-name"); } else if (!window.isActive() && now.isAfter(window.end())) { // Start the topology if it's stopped commandExecutor.startTopology("your-topology-name", "/path/to/your-topology.jar", "com.your.package.TopologyMain"); } } } // Helper to execute Storm CLI commands @Component public class TopologyCommandExecutor { public void stopTopology(String topologyName) { executeCommand("storm kill " + topologyName + " -w 10"); // Wait 10s for graceful shutdown } public void startTopology(String topologyName, String jarPath, String mainClass) { executeCommand("storm jar " + jarPath + " " + mainClass + " " + topologyName); } private void executeCommand(String command) { try { new ProcessBuilder(command.split(" ")) .inheritIO() .start() .waitFor(); } catch (IOException | InterruptedException e) { // Add proper logging and error handling here e.printStackTrace(); } } }
If you want to avoid shell commands (and make integration cleaner with Spring), use Storm’s built-in REST API (available in Storm 1.0+). This lets you manage topologies directly via HTTP calls:
Step 1: Confirm Storm REST API is enabled
The REST API is enabled by default in most Storm clusters, with a base URL likehttp://<nimbus-host>:8744/api/v1/.Step 2: Interact with the API from Spring Boot
UseWebClientorRestTemplateto send requests to Storm’s endpoints:@Service public class StormRestClient { private final WebClient webClient; public StormRestClient(WebClient.Builder webClientBuilder) { this.webClient = webClientBuilder.baseUrl("http://nimbus-host:8744/api/v1/").build(); } // Check if topology is active public boolean isTopologyRunning(String topologyName) { return webClient.get() .uri("/topology/summary") .retrieve() .bodyToMono(TopologySummaryResponse.class) .map(response -> response.topologies().stream() .anyMatch(t -> t.name().equals(topologyName) && t.status().equals("ACTIVE"))) .block(); } // Kill a running topology public void killTopology(String topologyName) { webClient.post() .uri("/topology/{topologyName}/kill", topologyName) .bodyValue(Map.of("waitSeconds", 10)) .retrieve() .toBodilessEntity() .block(); } // Note: Submitting topologies via REST requires jar uploads, which can be complex. // For simplicity, many teams use the CLI for submission and REST for kill/status checks. }Update your
TopologySchedulerServiceto use this client instead of shell commands for more maintainable code.
Don’t skip these to ensure reliability:
- Verify topology state first: Always check if the topology is already running/stopped before sending commands—prevents redundant operations or errors.
- Handle dynamic window updates: If your maintenance service can change windows mid-cycle, ensure your scheduler fetches the latest data on every check.
- Log all actions: Add detailed logging for every topology state change and maintenance window check—this makes debugging issues a breeze.
- Test in staging first: Validate the full flow (stop during maintenance, start after) with a non-production topology to avoid disrupting live data processing.
内容的提问来源于stack exchange,提问作者Alejandro Agapito Bautista

