Skip to content
Open
6 changes: 6 additions & 0 deletions adapter/runtime/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,12 @@
<version>4.0.3</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.apache.rocketmq.eventbridge</groupId>
<artifactId>rocketmq-eventbridge-metrics</artifactId>
<version>1.0.0</version>
<scope>compile</scope>
</dependency>
</dependencies>


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@

package org.apache.rocketmq.eventbridge.adapter.runtime;

import org.apache.rocketmq.eventbridge.BridgeMetricsManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.EventBusListener;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.EventMonitor;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.EventRuleTransfer;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.EventTargetTrigger;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.common.CirculatorContext;
Expand Down Expand Up @@ -69,9 +71,13 @@ public void initAndStart() throws Exception {
circulatorContext.initCirculatorContext(runnerConfigObserver.getTargetRunnerConfig());
runnerConfigObserver.registerListener(circulatorContext);
runnerConfigObserver.registerListener(eventSubscriber);
EventBusListener eventBusListener = new EventBusListener(circulatorContext, eventSubscriber, errorHandler);
EventRuleTransfer eventRuleTransfer = new EventRuleTransfer(circulatorContext, offsetManager, errorHandler);
EventTargetTrigger eventTargetPusher = new EventTargetTrigger(circulatorContext, offsetManager, errorHandler);

EventMonitor eventMonitor = new EventMonitor(eventSubscriber);
BridgeMetricsManager metricsManager = eventSubscriber.getMetricsManager();
EventBusListener eventBusListener = new EventBusListener(circulatorContext, eventSubscriber, errorHandler, metricsManager);
EventRuleTransfer eventRuleTransfer = new EventRuleTransfer(circulatorContext, offsetManager, errorHandler, metricsManager);
EventTargetTrigger eventTargetPusher = new EventTargetTrigger(circulatorContext, offsetManager, errorHandler, metricsManager);
RUNTIME_START_AND_SHUTDOWN.appendStartAndShutdown(eventMonitor);
RUNTIME_START_AND_SHUTDOWN.appendStartAndShutdown(eventBusListener);
RUNTIME_START_AND_SHUTDOWN.appendStartAndShutdown(eventRuleTransfer);
RUNTIME_START_AND_SHUTDOWN.appendStartAndShutdown(eventTargetPusher);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,13 @@

import com.google.common.collect.Lists;
import io.openmessaging.connector.api.data.ConnectRecord;

import java.util.ArrayList;
import java.util.List;
import java.util.Optional;

import org.apache.commons.collections.CollectionUtils;
import org.apache.rocketmq.eventbridge.BridgeMetricsManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.common.CirculatorContext;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.listener.EventSubscriber;
import org.apache.rocketmq.eventbridge.adapter.runtime.common.ServiceThread;
Expand All @@ -40,20 +45,23 @@ public class EventBusListener extends ServiceThread {
private final CirculatorContext circulatorContext;
private final EventSubscriber eventSubscriber;
private final ErrorHandler errorHandler;
private BridgeMetricsManager metricsManager;

public EventBusListener(CirculatorContext circulatorContext, EventSubscriber eventSubscriber,
ErrorHandler errorHandler) {
ErrorHandler errorHandler, BridgeMetricsManager metricsManager) {
this.circulatorContext = circulatorContext;
this.eventSubscriber = eventSubscriber;
this.errorHandler = errorHandler;
this.metricsManager = metricsManager;
}

@Override
public void run() {
while (!stopped) {
List<ConnectRecord> pullRecordList = Lists.newArrayList();
try {
pullRecordList = eventSubscriber.pull();
pullRecordList = Optional.ofNullable(eventSubscriber.pull()).orElse(new ArrayList<>());
BridgeMetricsManager.messagesInTotal.add(pullRecordList.size());
if (CollectionUtils.isEmpty(pullRecordList)) {
this.waitForRunning(1000);
continue;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.rocketmq.eventbridge.adapter.runtime.boot;

import org.apache.rocketmq.eventbridge.BridgeMetricsManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.listener.EventSubscriber;
import org.apache.rocketmq.eventbridge.adapter.runtime.common.ServiceThread;

public class EventMonitor extends ServiceThread {

private BridgeMetricsManager bridgeMetricsManager;

public EventMonitor(EventSubscriber eventSubscriber) {
this.bridgeMetricsManager = eventSubscriber.getMetricsManager();
}
@Override
public String getServiceName() {
return EventMonitor.class.getSimpleName();
}

@Override
public void run() {
bridgeMetricsManager.init();
}

@Override
public void shutdown() {
bridgeMetricsManager.shutdown();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import java.util.concurrent.CompletableFuture;
import javax.annotation.PostConstruct;
import org.apache.commons.collections.MapUtils;
import org.apache.rocketmq.eventbridge.BridgeMetricsManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.common.CirculatorContext;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.common.OffsetManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.transfer.TransformEngine;
Expand All @@ -47,12 +48,14 @@ public class EventRuleTransfer extends ServiceThread {
private final CirculatorContext circulatorContext;
private final OffsetManager offsetManager;
private final ErrorHandler errorHandler;
private BridgeMetricsManager metricsManager;

public EventRuleTransfer(CirculatorContext circulatorContext, OffsetManager offsetManager,
ErrorHandler errorHandler) {
ErrorHandler errorHandler, BridgeMetricsManager metricsManager) {
this.circulatorContext = circulatorContext;
this.offsetManager = offsetManager;
this.errorHandler = errorHandler;
this.metricsManager = metricsManager;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import java.util.concurrent.ExecutorService;

import org.apache.commons.collections.MapUtils;
import org.apache.rocketmq.eventbridge.BridgeMetricsManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.common.OffsetManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.common.CirculatorContext;
import org.apache.rocketmq.eventbridge.adapter.runtime.common.ServiceThread;
Expand All @@ -47,12 +48,14 @@ public class EventTargetTrigger extends ServiceThread {
private final OffsetManager offsetManager;
private final ErrorHandler errorHandler;
private volatile Integer batchSize = 100;
private BridgeMetricsManager metricsManager;

public EventTargetTrigger(CirculatorContext circulatorContext, OffsetManager offsetManager,
ErrorHandler errorHandler) {
ErrorHandler errorHandler, BridgeMetricsManager metricsManager) {
this.circulatorContext = circulatorContext;
this.offsetManager = offsetManager;
this.errorHandler = errorHandler;
this.metricsManager = metricsManager;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.rocketmq.eventbridge.adapter.runtime.boot.listener;

import io.openmessaging.connector.api.data.ConnectRecord;
import org.apache.rocketmq.eventbridge.BridgeMetricsManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.common.TargetRunnerListener;
import org.apache.rocketmq.eventbridge.adapter.runtime.common.entity.SubscribeRunnerKeys;
import org.apache.rocketmq.eventbridge.adapter.runtime.common.entity.TargetRunnerConfig;
Expand All @@ -34,6 +35,13 @@ public abstract class EventSubscriber implements TargetRunnerListener {
*/
public abstract void refresh(SubscribeRunnerKeys subscribeRunnerKeys, RefreshTypeEnum refreshTypeEnum);


/**
* fetch metrics configuration
* @return
*/
public abstract BridgeMetricsManager getMetricsManager();

/**
* Pull connect records from store, Blocking method when is empty.
*
Expand Down
12 changes: 8 additions & 4 deletions adapter/runtime/src/main/resources/runtime.properties
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,14 @@ rocketmq.namesrvAddr=localhost:9876
rocketmq.consumer.pullTimeOut = 3000
rocketmq.consumer.pullBatchSize=20
rocketmq.cluster.name=DefaultCluster
rocketmq.consumerGroup=default
## runtime
rumtimer.name=eventbridge-runtimer
runtimer.pluginpath=/Users/Local/eventbridge/plugin
runtimer.storePathRootDir=/Users/Local/eventbridge/store
rumtime.name=eventbridge-runtimer
## listener
listener.eventQueue.threshold=50000
listener.targetQueue.threshold=50000
listener.targetQueue.threshold=50000

## monitor
metrics.endpoint.host=127.0.0.1
metrics.endpoint.port=19090
metrics.collector.mode=2
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,9 @@
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.utils.NetworkUtil;
import org.apache.rocketmq.eventbridge.BridgeConfig;
import org.apache.rocketmq.eventbridge.BridgeMetricsManager;
import org.apache.rocketmq.eventbridge.adapter.runtime.boot.listener.EventSubscriber;
import org.apache.rocketmq.eventbridge.adapter.runtime.common.ServiceThread;
import org.apache.rocketmq.eventbridge.adapter.runtime.common.entity.SubscribeRunnerKeys;
Expand Down Expand Up @@ -85,6 +88,7 @@ public class RocketMQEventSubscriber extends EventSubscriber {
private Integer pullTimeOut;
private Integer pullBatchSize;

private BridgeConfig bridgeConfig;
private ClientConfig clientConfig;
private SessionCredentials sessionCredentials;
private String socksProxy;
Expand Down Expand Up @@ -118,6 +122,13 @@ public void refresh(SubscribeRunnerKeys subscribeRunnerKeys, RefreshTypeEnum ref
}
}

@Override
public BridgeMetricsManager getMetricsManager() {
BridgeMetricsManager metricsManager = new BridgeMetricsManager(bridgeConfig);
return metricsManager;
}


@Override
public List<ConnectRecord> pull() {
ArrayList<MessageExt> messages = new ArrayList<>();
Expand Down Expand Up @@ -185,11 +196,19 @@ private void initMqProperties() {
String socks5Password = properties.getProperty("rocketmq.consumer.socks5Password");
String socks5Endpoint = properties.getProperty("rocketmq.consumer.socks5Endpoint");

String metricsPromExporterHost = properties.getProperty("metrics.endpoint.host");
String metricsPromExporterPort = properties.getProperty("metrics.endpoint.port");
String metricsCollectorMode = properties.getProperty("metrics.collector.mode");
clientConfig.setNameSrvAddr(namesrvAddr);
clientConfig.setAccessChannel(AccessChannel.CLOUD.name().equals(accessChannel) ?
AccessChannel.CLOUD : AccessChannel.LOCAL);
clientConfig.setNamespace(namespace);
BridgeConfig bridgeConfig = new BridgeConfig();
bridgeConfig.setMetricsPromExporterHost(metricsPromExporterHost);
bridgeConfig.setMetricsPromExporterPort(Integer.parseInt(metricsPromExporterPort));
Comment thread
kaori-seasons marked this conversation as resolved.
Outdated
bridgeConfig.setMetricsExporterType(Integer.parseInt(metricsCollectorMode));
this.clientConfig = clientConfig;
this.bridgeConfig = bridgeConfig;

if (StringUtils.isNotBlank(accessKey) && StringUtils.isNotBlank(secretKey)) {
this.sessionCredentials = new SessionCredentials(accessKey, secretKey);
Expand Down
22 changes: 22 additions & 0 deletions common/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<opentelemetry.version>1.19.0</opentelemetry.version>
<opentelemetry-exporter-prometheus.version>1.19.0-alpha</opentelemetry-exporter-prometheus.version>
</properties>

<dependencies>
Expand Down Expand Up @@ -68,5 +70,25 @@
<artifactId>assertj-core</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.opentelemetry</groupId>
<artifactId>opentelemetry-exporter-otlp</artifactId>
<version>${opentelemetry.version}</version>
</dependency>
<dependency>
<groupId>io.opentelemetry</groupId>
<artifactId>opentelemetry-exporter-prometheus</artifactId>
<version>${opentelemetry-exporter-prometheus.version}</version>
</dependency>
<dependency>
<groupId>io.opentelemetry</groupId>
<artifactId>opentelemetry-exporter-logging</artifactId>
<version>${opentelemetry.version}</version>
</dependency>
<dependency>
<groupId>io.opentelemetry</groupId>
<artifactId>opentelemetry-sdk</artifactId>
<version>${opentelemetry.version}</version>
</dependency>
</dependencies>
</project>
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.rocketmq.eventbridge.metrics;

import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.metrics.LongCounter;
import io.opentelemetry.context.Context;

public class NopLongCounter implements LongCounter {
@Override public void add(long l) {

}

@Override public void add(long l, Attributes attributes) {

}

@Override public void add(long l, Attributes attributes, Context context) {

}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.rocketmq.eventbridge.metrics;

import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.metrics.LongHistogram;
import io.opentelemetry.context.Context;

public class NopLongHistogram implements LongHistogram {
@Override public void record(long l) {

}

@Override public void record(long l, Attributes attributes) {

}

@Override public void record(long l, Attributes attributes, Context context) {

}
}
Loading