You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在指定时段停用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:

1. Spring Boot Scheduled Tasks + Storm CLI

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 @Scheduled annotation 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();
            }
        }
    }
    
2. Use Storm's REST API for Programmatic Control

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 like http://<nimbus-host>:8744/api/v1/.

  • Step 2: Interact with the API from Spring Boot
    Use WebClient or RestTemplate to 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 TopologySchedulerService to use this client instead of shell commands for more maintainable code.

3. Critical Safeguards to Avoid Edge Cases

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 03:59:12