代码之家  ›  专栏  ›  技术社区  ›  suman j

如何使用ZooKeeper在Spring集成中实现分布式轮询锁

  •  2
  • suman j  · 技术社区  · 7 年前

    Spring集成具有ZooKeeper支持,如中所述 https://docs.spring.io/spring-integration/reference/html/zookeeper.html 然而,这份文件是如此模糊。

    它建议在bean下面添加,但没有给出当节点被授予领导权时如何启动/停止轮询器的详细信息。

    @Bean
    public LeaderInitiatorFactoryBean leaderInitiator(CuratorFramework client) {
        return new LeaderInitiatorFactoryBean()
                    .setClient(client)
                    .setPath("/siTest/")
                    .setRole("cluster");
    }
    

    我们有没有关于如何使用zookeeper确保在集群中任何时候只运行一次以下轮询器的示例?

    @Component
    public class EventsPoller {
    
        public void pullEvents() {
            //pull events should be run by only one node in the cluster at any time
        }
    }
    
    1 回复  |  直到 7 年前
        1
  •  2
  •   Artem Bilan    7 年前

    这个 LeaderInitiator 散发出 OnGrantedEvent OnRevokedEvent ,当其成为领导者且其领导地位被撤销时。

    看见 https://docs.spring.io/spring-integration/reference/html/messaging-endpoints-chapter.html#endpoint-roles https://docs.spring.io/spring-integration/reference/html/messaging-endpoints-chapter.html#leadership-event-handling 有关这些事件处理以及它如何影响特定角色中的组件的更多信息。

    SmartLifecycleRoleController 章请随时就此事提出JIRA,欢迎您的贡献!

    使现代化

    @RunWith(SpringRunner.class)
    @DirtiesContext
    public class LeaderInitiatorFactoryBeanTests extends ZookeeperTestSupport {
    
        private static CuratorFramework client;
    
        @Autowired
        private PollableChannel stringsChannel;
    
        @BeforeClass
        public static void getClient() throws Exception {
            client = createNewClient();
        }
    
        @AfterClass
        public static void closeClient() {
            if (client != null) {
                client.close();
            }
        }
    
        @Test
        public void test() {
            assertNotNull(this.stringsChannel.receive(10_000));
        }
    
    
        @Configuration
        @EnableIntegration
        public static class Config {
    
            @Bean
            public LeaderInitiatorFactoryBean leaderInitiator(CuratorFramework client) {
                return new LeaderInitiatorFactoryBean()
                        .setClient(client)
                        .setPath("/siTest/")
                        .setRole("foo");
            }
    
            @Bean
            public CuratorFramework client() {
                return LeaderInitiatorFactoryBeanTests.client;
            }
    
            @Bean
            @InboundChannelAdapter(channel = "stringsChannel", autoStartup = "false", poller = @Poller(fixedDelay = "100"))
            @Role("foo")
            public Supplier<String> inboundChannelAdapter() {
                return () -> "foo";
            }
    
            @Bean
            public PollableChannel stringsChannel() {
                return new QueueChannel();
            }
    
        }
    
    }
    

    我有这样的日志:

    2018-12-14 10:12:33,542 DEBUG [Curator-LeaderSelector-0] [org.springframework.integration.support.SmartLifecycleRoleController] - Starting [leaderInitiatorFactoryBeanTests.Config.inboundChannelAdapter.inboundChannelAdapter] in role foo
    2018-12-14 10:12:33,578 DEBUG [Curator-LeaderSelector-0] [org.springframework.integration.support.SmartLifecycleRoleController] - Stopping [leaderInitiatorFactoryBeanTests.Config.inboundChannelAdapter.inboundChannelAdapter] in role foo