Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@
import com.navercorp.pinpoint.web.dao.AgentEventDao;
import com.navercorp.pinpoint.web.dao.hbase.HbaseApplicationIndexDao;
import com.navercorp.pinpoint.web.dao.hbase.HbaseMapResponseTimeDao;
import com.navercorp.pinpoint.web.dao.hbase.HbaseMapStatisticsCallerDao;
import com.navercorp.pinpoint.web.dao.hbase.HbaseMapStatisticsCallerCompactDao;
import com.navercorp.pinpoint.web.dao.stat.AgentStatDao;
import com.navercorp.pinpoint.web.dao.stat.FileDescriptorDao;
import com.navercorp.pinpoint.web.vo.Application;
Expand Down Expand Up @@ -64,7 +64,7 @@ public class DataCollectorFactory {

private final HbaseApplicationIndexDao hbaseApplicationIndexDao;

private final HbaseMapStatisticsCallerDao mapStatisticsCallerDao;
private final HbaseMapStatisticsCallerCompactDao mapStatisticsCallerDao;

public DataCollectorFactory(HbaseMapResponseTimeDao hbaseMapResponseTimeDao,
@Qualifier("jvmGcDaoFactory") AgentStatDao<JvmGcBo> jvmGcDao,
Expand All @@ -73,7 +73,7 @@ public DataCollectorFactory(HbaseMapResponseTimeDao hbaseMapResponseTimeDao,
@Qualifier("fileDescriptorDaoFactory") FileDescriptorDao fileDescriptorDao,
AgentEventDao agentEventDao,
HbaseApplicationIndexDao hbaseApplicationIndexDao,
HbaseMapStatisticsCallerDao mapStatisticsCallerDao) {
HbaseMapStatisticsCallerCompactDao mapStatisticsCallerDao) {
this.hbaseMapResponseTimeDao = Objects.requireNonNull(hbaseMapResponseTimeDao, "hbaseMapResponseTimeDao");
this.jvmGcDao = Objects.requireNonNull(jvmGcDao, "jvmGcDao");
this.cpuLoadDao = Objects.requireNonNull(cpuLoadDao, "cpuLoadDao");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkCallDataMap;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkData;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkDataMap;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerDao;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerCompactDao;
import com.navercorp.pinpoint.web.vo.Application;
import com.navercorp.pinpoint.web.vo.Range;

Expand All @@ -36,13 +36,13 @@
public class MapStatisticsCallerDataCollector extends DataCollector {

private final Application application;
private final MapStatisticsCallerDao mapStatisticsCallerDao;
private final MapStatisticsCallerCompactDao mapStatisticsCallerDao;
private final long timeSlotEndTime;
private final long slotInterval;
private final Map<String, LinkCallData> calleeStatMap = new HashMap<>();
private final AtomicBoolean init = new AtomicBoolean(false); // need to consider a trace condition when checkers start simultaneously.

public MapStatisticsCallerDataCollector(DataCollectorCategory category, Application application, MapStatisticsCallerDao mapStatisticsCallerDao, long timeSlotEndTime, long slotInterval) {
public MapStatisticsCallerDataCollector(DataCollectorCategory category, Application application, MapStatisticsCallerCompactDao mapStatisticsCallerDao, long timeSlotEndTime, long slotInterval) {
super(category);
this.application = application;
this.mapStatisticsCallerDao = mapStatisticsCallerDao;
Expand Down
20 changes: 20 additions & 0 deletions batch/src/main/resources/applicationContext-batch-hbase.xml
Original file line number Diff line number Diff line change
Expand Up @@ -177,6 +177,26 @@
<constructor-arg type="int" value="8"/>
</bean>

<bean id="statisticsCalleeCompactRowKeyDistributor" class="com.sematext.hbase.wd.RowKeyDistributorByHashPrefix">
<constructor-arg ref="statisticsCalleeCompactHasher"/>
</bean>

<bean id="statisticsCalleeCompactHasher" class="com.navercorp.pinpoint.common.hbase.distributor.RangeOneByteSimpleHash">
<constructor-arg type="int" value="0"/>
<constructor-arg type="int" value="36"/>
<constructor-arg type="int" value="32"/>
</bean>

<bean id="statisticsCallerCompactRowKeyDistributor" class="com.sematext.hbase.wd.RowKeyDistributorByHashPrefix">
<constructor-arg ref="statisticsCallerCompactHasher"/>
</bean>

<bean id="statisticsCallerCompactHasher" class="com.navercorp.pinpoint.common.hbase.distributor.RangeOneByteSimpleHash">
<constructor-arg type="int" value="0"/>
<constructor-arg type="int" value="36"/>
<constructor-arg type="int" value="32"/>
</bean>

<bean id="slf4jCommonLoggerFactory" class="com.navercorp.pinpoint.common.server.util.Slf4jCommonLoggerFactory">
</bean>
<bean id="typeLoaderService" class="com.navercorp.pinpoint.common.server.util.ServerTraceMetadataLoaderService">
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkCallDataMap;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkData;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkDataMap;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerDao;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerCompactDao;
import com.navercorp.pinpoint.web.vo.Application;
import com.navercorp.pinpoint.web.vo.Range;
import org.junit.BeforeClass;
Expand All @@ -46,11 +46,11 @@ public class ErrorCountToCalleCheckerTest {
private static final String FROM_SERVICE_NAME = "from_local_service";
private static final String TO_SERVICE_NAME = "to_local_service";
private static final String SERVICE_TYPE = "tomcat";
public static MapStatisticsCallerDao dao;
public static MapStatisticsCallerCompactDao dao;

@BeforeClass
public static void before() {
dao = new MapStatisticsCallerDao() {
dao = new MapStatisticsCallerCompactDao() {

@Override
public LinkDataMap selectCaller(Application callerApplication, Range range) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

package com.navercorp.pinpoint.batch.alarm.checker;


import com.navercorp.pinpoint.web.alarm.CheckerCategory;
import com.navercorp.pinpoint.web.alarm.DataCollectorCategory;
import com.navercorp.pinpoint.batch.alarm.collector.MapStatisticsCallerDataCollector;
Expand All @@ -25,7 +26,7 @@
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkCallDataMap;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkData;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkDataMap;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerDao;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerCompactDao;
import com.navercorp.pinpoint.web.vo.Application;
import com.navercorp.pinpoint.web.vo.Range;
import org.junit.BeforeClass;
Expand All @@ -43,11 +44,11 @@ public class ErrorRateToCalleCheckerTest {
private static final String TO_SERVICE_NAME = "to_local_service";
private static final String SERVICE_TYPE = "tomcat";

public static MapStatisticsCallerDao dao;
public static MapStatisticsCallerCompactDao dao;

@BeforeClass
public static void before() {
dao = new MapStatisticsCallerDao() {
dao = new MapStatisticsCallerCompactDao() {

@Override
public LinkDataMap selectCaller(Application callerApplication, Range range) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,14 +18,15 @@

import com.navercorp.pinpoint.web.alarm.DataCollectorCategory;
import com.navercorp.pinpoint.common.trace.ServiceType;

import com.navercorp.pinpoint.web.alarm.CheckerCategory;
import com.navercorp.pinpoint.batch.alarm.collector.MapStatisticsCallerDataCollector;
import com.navercorp.pinpoint.web.alarm.vo.Rule;
import com.navercorp.pinpoint.web.applicationmap.histogram.TimeHistogram;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkCallDataMap;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkData;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkDataMap;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerDao;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerCompactDao;
import com.navercorp.pinpoint.web.vo.Application;
import com.navercorp.pinpoint.web.vo.Range;
import org.junit.BeforeClass;
Expand All @@ -42,11 +43,11 @@ public class SlowCountToCalleCheckerTest {
private static final String FROM_SERVICE_NAME = "from_local_service";
private static final String TO_SERVICE_NAME = "to_local_service";
private static final String SERVICE_TYPE = "tomcat";
public static MapStatisticsCallerDao dao;
public static MapStatisticsCallerCompactDao dao;

@BeforeClass
public static void before() {
dao = new MapStatisticsCallerDao() {
dao = new MapStatisticsCallerCompactDao() {

@Override
public LinkDataMap selectCaller(Application callerApplication, Range range) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

package com.navercorp.pinpoint.batch.alarm.checker;


import com.navercorp.pinpoint.web.alarm.CheckerCategory;
import com.navercorp.pinpoint.web.alarm.DataCollectorCategory;
import com.navercorp.pinpoint.batch.alarm.collector.MapStatisticsCallerDataCollector;
Expand All @@ -25,7 +26,7 @@
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkCallDataMap;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkData;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkDataMap;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerDao;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerCompactDao;
import com.navercorp.pinpoint.web.vo.Application;
import com.navercorp.pinpoint.web.vo.Range;
import org.junit.BeforeClass;
Expand All @@ -42,11 +43,11 @@ public class SlowRateToCalleCheckerTest {
private static final String FROM_SERVICE_NAME = "from_local_service";
private static final String TO_SERVICE_NAME = "to_local_service";
private static final String SERVICE_TYPE = "tomcat";
public static MapStatisticsCallerDao dao;
public static MapStatisticsCallerCompactDao dao;

@BeforeClass
public static void before() {
dao = new MapStatisticsCallerDao() {
dao = new MapStatisticsCallerCompactDao() {

@Override
public LinkDataMap selectCaller(Application callerApplication, Range range) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkCallDataMap;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkData;
import com.navercorp.pinpoint.web.applicationmap.rawdata.LinkDataMap;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerDao;
import com.navercorp.pinpoint.web.dao.MapStatisticsCallerCompactDao;
import com.navercorp.pinpoint.web.vo.Application;
import com.navercorp.pinpoint.web.vo.Range;
import org.junit.BeforeClass;
Expand All @@ -42,11 +42,11 @@ public class TotalCountToCalleeCheckerTest {
private static final String FROM_SERVICE_NAME = "from_local_service";
private static final String TO_SERVICE_NAME = "to_local_service";
private static final String SERVICE_TYPE = "tomcat";
public static MapStatisticsCallerDao dao;
public static MapStatisticsCallerCompactDao dao;

@BeforeClass
public static void before() {
dao = new MapStatisticsCallerDao() {
dao = new MapStatisticsCallerCompactDao() {

@Override
public LinkDataMap selectCaller(Application callerApplication, Range range) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
/*
* Copyright 2021 NAVER Corp.
*
* Licensed 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 com.navercorp.pinpoint.collector.dao;

import com.navercorp.pinpoint.common.trace.ServiceType;

public interface MapStatisticsCalleeCompactDao extends CachedStatisticsDao {

void update(String calleeApplicationName, ServiceType calleeServiceType, String callerApplicationName, ServiceType callerServiceType, int elapsed, boolean isError);

}
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
/*
* Copyright 2021 NAVER Corp.
*
* Licensed 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 com.navercorp.pinpoint.collector.dao;

import com.navercorp.pinpoint.common.trace.ServiceType;

public interface MapStatisticsCallerCompactDao extends CachedStatisticsDao {
void update(String callerApplicationName, ServiceType callerServiceType, String calleeApplicationName, ServiceType calleeServiceType, String calleeHost, int elapsed, boolean isError);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
/*
* Copyright 2021 NAVER Corp.
*
* Licensed 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 com.navercorp.pinpoint.collector.dao.hbase;

import com.navercorp.pinpoint.collector.dao.MapStatisticsCalleeCompactDao;
import com.navercorp.pinpoint.collector.dao.hbase.statistics.BulkWriter;
import com.navercorp.pinpoint.collector.dao.hbase.statistics.CallRowKey;
import com.navercorp.pinpoint.collector.dao.hbase.statistics.CallerCompactColumnName;
import com.navercorp.pinpoint.collector.dao.hbase.statistics.ColumnName;
import com.navercorp.pinpoint.collector.dao.hbase.statistics.MapLinkConfiguration;
import com.navercorp.pinpoint.collector.dao.hbase.statistics.RowKey;
import com.navercorp.pinpoint.common.server.util.AcceptedTimeService;
import com.navercorp.pinpoint.common.server.util.ApplicationMapStatisticsUtils;
import com.navercorp.pinpoint.common.server.util.TimeSlot;
import com.navercorp.pinpoint.common.trace.HistogramSchema;
import com.navercorp.pinpoint.common.trace.ServiceType;
import org.apache.logging.log4j.Logger;
import org.apache.logging.log4j.LogManager;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Repository;

import java.util.Objects;

/**
* Update statistics of callee node
*/
@Repository
public class HbaseMapStatisticsCalleeCompactDao implements MapStatisticsCalleeCompactDao {

private final Logger logger = LogManager.getLogger(this.getClass());

private final AcceptedTimeService acceptedTimeService;

private final TimeSlot timeSlot;

private final IgnoreStatFilter ignoreStatFilter;
private final BulkWriter bulkWriter;
private final MapLinkConfiguration mapLinkConfiguration;

@Autowired
public HbaseMapStatisticsCalleeCompactDao(MapLinkConfiguration mapLinkConfiguration,
IgnoreStatFilter ignoreStatFilter,
AcceptedTimeService acceptedTimeService, TimeSlot timeSlot,
@Qualifier("calleeCompactBulkWriter") BulkWriter bulkWriter) {
this.mapLinkConfiguration = Objects.requireNonNull(mapLinkConfiguration, "mapLinkConfiguration");
this.ignoreStatFilter = Objects.requireNonNull(ignoreStatFilter, "ignoreStatFilter");
this.acceptedTimeService = Objects.requireNonNull(acceptedTimeService, "acceptedTimeService");
this.timeSlot = Objects.requireNonNull(timeSlot, "timeSlot");

this.bulkWriter = Objects.requireNonNull(bulkWriter, "bulkWriter");
}

@Override
public void update(String calleeApplicationName, ServiceType calleeServiceType, String callerApplicationName, ServiceType callerServiceType, int elapsed, boolean isError) {
Objects.requireNonNull(calleeApplicationName, "calleeApplicationName");
Objects.requireNonNull(callerApplicationName, "callerApplicationName");

if (logger.isDebugEnabled()) {
logger.debug("[Callee] {} ({}) <- {} ({})",
calleeApplicationName, calleeServiceType, callerApplicationName, callerServiceType);
}

// make row key. rowkey is me
final long acceptedTime = acceptedTimeService.getAcceptedTime();
final long rowTimeSlot = timeSlot.getTimeSlot(acceptedTime);
final RowKey calleeRowKey = new CallRowKey(calleeApplicationName, calleeServiceType.getCode(), rowTimeSlot);
final short callerSlotNumber = ApplicationMapStatisticsUtils.getSlotNumber(calleeServiceType, elapsed, isError);
HistogramSchema histogramSchema = calleeServiceType.getHistogramSchema();

final ColumnName callerColumnName = new CallerCompactColumnName(callerServiceType.getCode(), callerApplicationName, callerSlotNumber);
this.bulkWriter.increment(calleeRowKey, callerColumnName);

if (mapLinkConfiguration.isEnableAvg()) {
final ColumnName sumColumnName = new CallerCompactColumnName(callerServiceType.getCode(), callerApplicationName, histogramSchema.getSumStatSlot().getSlotTime());
this.bulkWriter.increment(calleeRowKey, sumColumnName, elapsed);
}
if (mapLinkConfiguration.isEnableMax()) {
final ColumnName maxColumnName = new CallerCompactColumnName(callerServiceType.getCode(), callerApplicationName, histogramSchema.getMaxStatSlot().getSlotTime());
this.bulkWriter.updateMax(calleeRowKey, maxColumnName, elapsed);
}
}

@Override
public void flushLink() {
this.bulkWriter.flushLink();
}

@Override
public void flushAvgMax() {
this.bulkWriter.flushAvgMax();
}
}
Loading