• 设为首页
  • 点击收藏
  • 手机版
    手机扫一扫访问
    迪恩网络手机版
  • 关注官方公众号
    微信扫一扫关注
    迪恩网络公众号

Java ConnectorContext类代码示例

原作者: [db:作者] 来自: [db:来源] 收藏 邀请

本文整理汇总了Java中org.apache.kafka.connect.connector.ConnectorContext的典型用法代码示例。如果您正苦于以下问题:Java ConnectorContext类的具体用法?Java ConnectorContext怎么用?Java ConnectorContext使用的例子?那么恭喜您, 这里精选的类代码示例或许可以为您提供帮助。



ConnectorContext类属于org.apache.kafka.connect.connector包,在下文中一共展示了ConnectorContext类的20个代码示例,这些例子默认根据受欢迎程度排序。您可以为喜欢或者感觉有用的代码点赞,您的评价将有助于我们的系统推荐出更棒的Java代码示例。

示例1: initialize

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
public void initialize(ConnectorConfig connectorConfig) {
    try {
        this.config = connectorConfig.originalsStrings();
        log.debug("Initializing connector {} with config {}", connName, config);

        connector.initialize(new ConnectorContext() {
            @Override
            public void requestTaskReconfiguration() {
                ctx.requestTaskReconfiguration();
            }

            @Override
            public void raiseError(Exception e) {
                log.error("Connector raised an error {}", connName, e);
                onFailure(e);
                ctx.raiseError(e);
            }
        });
    } catch (Throwable t) {
        log.error("Error initializing connector {}", connName, t);
        onFailure(t);
    }
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:24,代码来源:WorkerConnector.java


示例2: testStartupPaused

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testStartupPaused() {
    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    // connector never gets started

    listener.onPause(CONNECTOR);
    expectLastCall();

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.PAUSED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:24,代码来源:WorkerConnectorTest.java


示例3: testStartupFailure

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testStartupFailure() {
    RuntimeException exception = new RuntimeException();

    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    connector.start(CONFIG);
    expectLastCall().andThrow(exception);

    listener.onFailure(CONNECTOR, exception);
    expectLastCall();

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:27,代码来源:WorkerConnectorTest.java


示例4: testJoinAssignment

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testJoinAssignment() throws Exception {
    // Join group and get assignment
    EasyMock.expect(member.memberId()).andStubReturn("member");
    EasyMock.expect(worker.getPlugins()).andReturn(plugins);
    expectRebalance(1, Arrays.asList(CONN1), Arrays.asList(TASK1));
    expectPostRebalanceCatchup(SNAPSHOT);
    worker.startConnector(EasyMock.eq(CONN1), EasyMock.<Map<String, String>>anyObject(), EasyMock.<ConnectorContext>anyObject(),
            EasyMock.eq(herder), EasyMock.eq(TargetState.STARTED));
    PowerMock.expectLastCall().andReturn(true);
    EasyMock.expect(worker.isRunning(CONN1)).andReturn(true);

    EasyMock.expect(worker.connectorTaskConfigs(CONN1, MAX_TASKS, null)).andReturn(TASK_CONFIGS);
    worker.startTask(EasyMock.eq(TASK1), EasyMock.<Map<String, String>>anyObject(), EasyMock.<Map<String, String>>anyObject(),
            EasyMock.eq(herder), EasyMock.eq(TargetState.STARTED));
    PowerMock.expectLastCall().andReturn(true);
    member.poll(EasyMock.anyInt());
    PowerMock.expectLastCall();

    PowerMock.replayAll();

    herder.tick();

    PowerMock.verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:26,代码来源:DistributedHerderTest.java


示例5: testJoinLeaderCatchUpFails

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testJoinLeaderCatchUpFails() throws Exception {
    // Join group and as leader fail to do assignment
    EasyMock.expect(member.memberId()).andStubReturn("leader");
    expectRebalance(Collections.<String>emptyList(), Collections.<ConnectorTaskId>emptyList(),
            ConnectProtocol.Assignment.CONFIG_MISMATCH, 1, Collections.<String>emptyList(),
            Collections.<ConnectorTaskId>emptyList());
    // Reading to end of log times out
    configBackingStore.refresh(EasyMock.anyLong(), EasyMock.anyObject(TimeUnit.class));
    EasyMock.expectLastCall().andThrow(new TimeoutException());
    member.maybeLeaveGroup();
    EasyMock.expectLastCall();
    PowerMock.expectPrivate(herder, "backoff", DistributedConfig.WORKER_UNSYNC_BACKOFF_MS_DEFAULT);
    member.requestRejoin();

    // After backoff, restart the process and this time succeed
    expectRebalance(1, Arrays.asList(CONN1), Arrays.asList(TASK1));
    expectPostRebalanceCatchup(SNAPSHOT);

    worker.startConnector(EasyMock.eq(CONN1), EasyMock.<Map<String, String>>anyObject(), EasyMock.<ConnectorContext>anyObject(),
            EasyMock.eq(herder), EasyMock.eq(TargetState.STARTED));
    PowerMock.expectLastCall().andReturn(true);
    EasyMock.expect(worker.getPlugins()).andReturn(plugins);
    EasyMock.expect(worker.connectorTaskConfigs(CONN1, MAX_TASKS, null)).andReturn(TASK_CONFIGS);
    worker.startTask(EasyMock.eq(TASK1), EasyMock.<Map<String, String>>anyObject(), EasyMock.<Map<String, String>>anyObject(),
            EasyMock.eq(herder), EasyMock.eq(TargetState.STARTED));
    PowerMock.expectLastCall().andReturn(true);
    EasyMock.expect(worker.isRunning(CONN1)).andReturn(true);
    member.poll(EasyMock.anyInt());
    PowerMock.expectLastCall();

    PowerMock.replayAll();

    herder.tick();
    herder.tick();

    PowerMock.verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:39,代码来源:DistributedHerderTest.java


示例6: PlogMonitorThread

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
/**
 * Create PLOG monitoring thread
 * 
 * @param context Connector context to use for monitoring thread to
 *                request task reconfiguration
 * @param fileManager the PLOG file manager to use for producing PLOGs
 * @param infoTopic the replicate info topic reader/writer
 * @param catalog the persistent replicate offset catalog read from topic
 * @param transactionTopicName name of transaction info topic to use
 * @param publishTransactions whether or not to publish transaction messages
 * 
 * @throws Exception when any startup error occurs
 */
public PlogMonitorThread (
    final ConnectorContext context, 
    final PlogFileManager fileManager,
    final ReplicateInfoTopic infoTopic,
    final Map <String, ReplicateInfo> catalog,
    final String transactionTopicName,
    final boolean publishTransactions
) throws Exception 
{
    this.context     = context;
    this.fileManager = fileManager;
    this.infoTopic   = infoTopic;
    this.catalog     = catalog;
    this.transactionTopicName = transactionTopicName;
    this.publishTransactions  = publishTransactions;
    this.shutdownLatch = new CountDownLatch(1);
    this.updateTasks   = new AtomicBoolean(false);
    this.startTasks    = new AtomicBoolean(catalog.isEmpty());
    
    setName ("PLOG monitor thread");
    
    /* add shutdown hook */
    shutdownHook = new ShutdownHook();
    Runtime.getRuntime().addShutdownHook(shutdownHook);
}
 
开发者ID:dbvisitsoftware,项目名称:replicate-connector-for-kafka,代码行数:39,代码来源:PlogMonitorThread.java


示例7: setup

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Before
public void setup() {
    connector = new FileStreamSinkConnector();
    ctx = PowerMock.createMock(ConnectorContext.class);
    connector.initialize(ctx);

    sinkProperties = new HashMap<>();
    sinkProperties.put(SinkConnector.TOPICS_CONFIG, MULTIPLE_TOPICS);
    sinkProperties.put(FileStreamSinkConnector.FILE_CONFIG, FILENAME);
}
 
开发者ID:wngn123,项目名称:wngn-jms-kafka,代码行数:11,代码来源:FileStreamSinkConnectorTest.java


示例8: setup

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Before
public void setup() {
    connector = new FileStreamSourceConnector();
    ctx = PowerMock.createMock(ConnectorContext.class);
    connector.initialize(ctx);

    sourceProperties = new HashMap<>();
    sourceProperties.put(FileStreamSourceConnector.TOPIC_CONFIG, SINGLE_TOPIC);
    sourceProperties.put(FileStreamSourceConnector.FILE_CONFIG, FILENAME);
}
 
开发者ID:wngn123,项目名称:wngn-jms-kafka,代码行数:11,代码来源:FileStreamSourceConnectorTest.java


示例9: setup

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Before
public void setup(){
  connector = new AmpoolSinkConnector();
  context = PowerMock.createMock(ConnectorContext.class);
  connector.initialize(context);

  sinkProperties = new HashMap<>();
  sinkProperties.put("locator.host", "localhost");
  sinkProperties.put("locator.port", Integer.toString(12345));
  sinkProperties.put("batch.size", Integer.toString(100));
  sinkProperties.put("topics", "topic1, topic2, topic3");
  sinkProperties.put("ampool.tables", "table1, table2, table3");
}
 
开发者ID:ampool,项目名称:monarch,代码行数:14,代码来源:AmpoolSinkConnectorTest.java


示例10: startConnector

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
/**
 * Start a connector managed by this worker.
 *
 * @param connName the connector name.
 * @param connProps the properties of the connector.
 * @param ctx the connector runtime context.
 * @param statusListener a listener for the runtime status transitions of the connector.
 * @param initialState the initial state of the connector.
 * @return true if the connector started successfully.
 */
public boolean startConnector(
        String connName,
        Map<String, String> connProps,
        ConnectorContext ctx,
        ConnectorStatus.Listener statusListener,
        TargetState initialState
) {
    if (connectors.containsKey(connName))
        throw new ConnectException("Connector with name " + connName + " already exists");

    final WorkerConnector workerConnector;
    ClassLoader savedLoader = plugins.currentThreadLoader();
    try {
        final ConnectorConfig connConfig = new ConnectorConfig(plugins, connProps);
        final String connClass = connConfig.getString(ConnectorConfig.CONNECTOR_CLASS_CONFIG);
        log.info("Creating connector {} of type {}", connName, connClass);
        final Connector connector = plugins.newConnector(connClass);
        workerConnector = new WorkerConnector(connName, connector, ctx, statusListener);
        log.info("Instantiated connector {} with version {} of type {}", connName, connector.version(), connector.getClass());
        savedLoader = plugins.compareAndSwapLoaders(connector);
        workerConnector.initialize(connConfig);
        workerConnector.transitionTo(initialState);
        Plugins.compareAndSwapLoaders(savedLoader);
    } catch (Throwable t) {
        log.error("Failed to start connector {}", connName, t);
        // Can't be put in a finally block because it needs to be swapped before the call on
        // statusListener
        Plugins.compareAndSwapLoaders(savedLoader);
        statusListener.onFailure(connName, t);
        return false;
    }

    WorkerConnector existing = connectors.putIfAbsent(connName, workerConnector);
    if (existing != null)
        throw new ConnectException("Connector with name " + connName + " already exists");

    log.info("Finished creating connector {}", connName);
    return true;
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:50,代码来源:Worker.java


示例11: WorkerConnector

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
public WorkerConnector(String connName,
                       Connector connector,
                       ConnectorContext ctx,
                       ConnectorStatus.Listener statusListener) {
    this.connName = connName;
    this.ctx = ctx;
    this.connector = connector;
    this.statusListener = statusListener;
    this.state = State.INIT;
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:11,代码来源:WorkerConnector.java


示例12: startConnector

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
private boolean startConnector(String connectorName) {
    log.info("Starting connector {}", connectorName);
    final Map<String, String> configProps = configState.connectorConfig(connectorName);
    final ConnectorContext ctx = new HerderConnectorContext(this, connectorName);
    final TargetState initialState = configState.targetState(connectorName);
    boolean started = worker.startConnector(connectorName, configProps, ctx, this, initialState);

    // Immediately request configuration since this could be a brand new connector. However, also only update those
    // task configs if they are actually different from the existing ones to avoid unnecessary updates when this is
    // just restoring an existing connector.
    if (started && initialState == TargetState.STARTED)
        reconfigureConnectorTasksWithRetry(connectorName);

    return started;
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:16,代码来源:DistributedHerder.java


示例13: testStartupAndShutdown

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testStartupAndShutdown() {
    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    connector.start(CONFIG);
    expectLastCall();

    listener.onStartup(CONNECTOR);
    expectLastCall();

    connector.stop();
    expectLastCall();

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:28,代码来源:WorkerConnectorTest.java


示例14: testStartupAndPause

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testStartupAndPause() {
    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    connector.start(CONFIG);
    expectLastCall();

    listener.onStartup(CONNECTOR);
    expectLastCall();

    connector.stop();
    expectLastCall();

    listener.onPause(CONNECTOR);
    expectLastCall();

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.transitionTo(TargetState.PAUSED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:32,代码来源:WorkerConnectorTest.java


示例15: testOnResume

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testOnResume() {
    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    listener.onPause(CONNECTOR);
    expectLastCall();

    connector.start(CONFIG);
    expectLastCall();

    listener.onResume(CONNECTOR);
    expectLastCall();

    connector.stop();
    expectLastCall();

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.PAUSED);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:32,代码来源:WorkerConnectorTest.java


示例16: testShutdownFailure

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testShutdownFailure() {
    RuntimeException exception = new RuntimeException();

    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    connector.start(CONFIG);
    expectLastCall();

    listener.onStartup(CONNECTOR);
    expectLastCall();

    connector.stop();
    expectLastCall().andThrow(exception);

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:30,代码来源:WorkerConnectorTest.java


示例17: testTransitionStartedToStarted

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testTransitionStartedToStarted() {
    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    connector.start(CONFIG);
    expectLastCall();

    // expect only one call to onStartup()
    listener.onStartup(CONNECTOR);
    expectLastCall();

    connector.stop();
    expectLastCall();

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:30,代码来源:WorkerConnectorTest.java


示例18: testTransitionPausedToPaused

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testTransitionPausedToPaused() {
    connector.initialize(EasyMock.notNull(ConnectorContext.class));
    expectLastCall();

    connector.start(CONFIG);
    expectLastCall();

    listener.onStartup(CONNECTOR);
    expectLastCall();

    connector.stop();
    expectLastCall();

    listener.onPause(CONNECTOR);
    expectLastCall();

    listener.onShutdown(CONNECTOR);
    expectLastCall();

    replayAll();

    WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, ctx, listener);

    workerConnector.initialize(connectorConfig);
    workerConnector.transitionTo(TargetState.STARTED);
    workerConnector.transitionTo(TargetState.PAUSED);
    workerConnector.transitionTo(TargetState.PAUSED);
    workerConnector.shutdown();

    verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:33,代码来源:WorkerConnectorTest.java


示例19: testStartConnectorFailure

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testStartConnectorFailure() throws Exception {
    expectConverters();
    expectStartStorage();

    ConnectorContext ctx = PowerMock.createMock(ConnectorContext.class);

    Map<String, String> props = new HashMap<>();
    props.put(SinkConnectorConfig.TOPICS_CONFIG, "foo,bar");
    props.put(ConnectorConfig.TASKS_MAX_CONFIG, "1");
    props.put(ConnectorConfig.NAME_CONFIG, CONNECTOR_ID);
    props.put(ConnectorConfig.CONNECTOR_CLASS_CONFIG, "java.util.HashMap"); // Bad connector class name

    EasyMock.expect(plugins.currentThreadLoader()).andReturn(delegatingLoader);
    EasyMock.expect(plugins.newConnector(EasyMock.anyString()))
            .andThrow(new ConnectException("Failed to find Connector"));

    EasyMock.expect(Plugins.compareAndSwapLoaders(delegatingLoader))
            .andReturn(pluginLoader);

    connectorStatusListener.onFailure(
            EasyMock.eq(CONNECTOR_ID),
            EasyMock.<ConnectException>anyObject()
    );
    EasyMock.expectLastCall();

    PowerMock.replayAll();

    worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore);
    worker.start();

    assertFalse(worker.startConnector(CONNECTOR_ID, props, ctx, connectorStatusListener, TargetState.STARTED));

    assertEquals(Collections.emptySet(), worker.connectorNames());

    assertFalse(worker.stopConnector(CONNECTOR_ID));

    PowerMock.verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:40,代码来源:WorkerTest.java


示例20: testRebalance

import org.apache.kafka.connect.connector.ConnectorContext; //导入依赖的package包/类
@Test
public void testRebalance() throws Exception {
    // Join group and get assignment
    EasyMock.expect(member.memberId()).andStubReturn("member");
    EasyMock.expect(worker.getPlugins()).andReturn(plugins);
    expectRebalance(1, Arrays.asList(CONN1), Arrays.asList(TASK1));
    expectPostRebalanceCatchup(SNAPSHOT);
    worker.startConnector(EasyMock.eq(CONN1), EasyMock.<Map<String, String>>anyObject(), EasyMock.<ConnectorContext>anyObject(),
            EasyMock.eq(herder), EasyMock.eq(TargetState.STARTED));
    PowerMock.expectLastCall().andReturn(true);
    EasyMock.expect(worker.isRunning(CONN1)).andReturn(true);
    EasyMock.expect(worker.connectorTaskConfigs(CONN1, MAX_TASKS, null)).andReturn(TASK_CONFIGS);
    worker.startTask(EasyMock.eq(TASK1), EasyMock.<Map<String, String>>anyObject(), EasyMock.<Map<String, String>>anyObject(),
            EasyMock.eq(herder), EasyMock.eq(TargetState.STARTED));
    PowerMock.expectLastCall().andReturn(true);
    member.poll(EasyMock.anyInt());
    PowerMock.expectLastCall();

    EasyMock.expect(worker.getPlugins()).andReturn(plugins);
    expectRebalance(Arrays.asList(CONN1), Arrays.asList(TASK1), ConnectProtocol.Assignment.NO_ERROR,
            1, Arrays.asList(CONN1), Arrays.<ConnectorTaskId>asList());

    // and the new assignment started
    worker.startConnector(EasyMock.eq(CONN1), EasyMock.<Map<String, String>>anyObject(), EasyMock.<ConnectorContext>anyObject(),
            EasyMock.eq(herder), EasyMock.eq(TargetState.STARTED));
    PowerMock.expectLastCall().andReturn(true);
    EasyMock.expect(worker.isRunning(CONN1)).andReturn(true);
    EasyMock.expect(worker.connectorTaskConfigs(CONN1, MAX_TASKS, null)).andReturn(TASK_CONFIGS);
    member.poll(EasyMock.anyInt());
    PowerMock.expectLastCall();

    PowerMock.replayAll();

    herder.tick();
    herder.tick();

    PowerMock.verifyAll();
}
 
开发者ID:YMCoding,项目名称:kafka-0.11.0.0-src-with-comment,代码行数:39,代码来源:DistributedHerderTest.java



注:本文中的org.apache.kafka.connect.connector.ConnectorContext类示例整理自Github/MSDocs等源码及文档管理平台,相关代码片段筛选自各路编程大神贡献的开源项目,源码版权归原作者所有,传播和使用请参考对应项目的License;未经允许,请勿转载。


鲜花

握手

雷人

路过

鸡蛋
该文章已有0人参与评论

请发表评论

全部评论

专题导读
上一篇:
Java TokenMetadata类代码示例发布时间:2022-05-22
下一篇:
Java YarnScheduler类代码示例发布时间:2022-05-22
热门推荐
阅读排行榜

扫描微信二维码

查看手机版网站

随时了解更新最新资讯

139-2527-9053

在线客服(服务时间 9:00~18:00)

在线QQ客服
地址:深圳市南山区西丽大学城创智工业园
电邮:jeky_zhao#qq.com
移动电话:139-2527-9053

Powered by 互联科技 X3.4© 2001-2213 极客世界.|Sitemap