diff --git a/apm-collector/apm-collector-component/client-component/pom.xml b/apm-collector/apm-collector-component/client-component/pom.xml index 612a3a585..21e367804 100644 --- a/apm-collector/apm-collector-component/client-component/pom.xml +++ b/apm-collector/apm-collector-component/client-component/pom.xml @@ -12,4 +12,47 @@ client-component jar + + + com.h2database + h2 + 1.4.196 + + + redis.clients + jedis + 2.9.0 + + + org.elasticsearch.client + transport + 5.5.0 + + + snakeyaml + org.yaml + + + + + org.apache.zookeeper + zookeeper + 3.4.10 + + + slf4j-api + org.slf4j + + + slf4j-log4j12 + org.slf4j + + + + + io.grpc + grpc-core + 1.4.0 + + \ No newline at end of file diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java new file mode 100644 index 000000000..791199a17 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java @@ -0,0 +1,164 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.elasticsearch; + +import java.net.InetAddress; +import java.net.UnknownHostException; +import java.util.LinkedList; +import java.util.List; +import java.util.concurrent.ExecutionException; +import org.elasticsearch.action.admin.indices.create.CreateIndexResponse; +import org.elasticsearch.action.admin.indices.delete.DeleteIndexResponse; +import org.elasticsearch.action.admin.indices.exists.indices.IndicesExistsResponse; +import org.elasticsearch.action.bulk.BulkRequestBuilder; +import org.elasticsearch.action.get.GetRequestBuilder; +import org.elasticsearch.action.get.MultiGetRequestBuilder; +import org.elasticsearch.action.index.IndexRequestBuilder; +import org.elasticsearch.action.search.SearchRequestBuilder; +import org.elasticsearch.action.update.UpdateRequest; +import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.elasticsearch.client.IndicesAdminClient; +import org.elasticsearch.common.settings.Settings; +import org.elasticsearch.common.transport.InetSocketTransportAddress; +import org.elasticsearch.common.xcontent.XContentBuilder; +import org.elasticsearch.transport.client.PreBuiltTransportClient; +import org.skywalking.apm.collector.core.client.Client; +import org.skywalking.apm.collector.core.client.ClientException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class ElasticSearchClient implements Client { + + private final Logger logger = LoggerFactory.getLogger(ElasticSearchClient.class); + + private org.elasticsearch.client.Client client; + + private final String clusterName; + + private final Boolean clusterTransportSniffer; + + private final String clusterNodes; + + public ElasticSearchClient(String clusterName, Boolean clusterTransportSniffer, String clusterNodes) { + this.clusterName = clusterName; + this.clusterTransportSniffer = clusterTransportSniffer; + this.clusterNodes = clusterNodes; + } + + @Override public void initialize() throws ClientException { + Settings settings = Settings.builder() + .put("cluster.name", clusterName) + .put("client.transport.sniff", clusterTransportSniffer) + .build(); + + client = new PreBuiltTransportClient(settings); + + List pairsList = parseClusterNodes(clusterNodes); + for (AddressPairs pairs : pairsList) { + try { + ((PreBuiltTransportClient)client).addTransportAddress(new InetSocketTransportAddress(InetAddress.getByName(pairs.host), pairs.port)); + } catch (UnknownHostException e) { + throw new ElasticSearchClientException(e.getMessage(), e); + } + } + } + + @Override public void shutdown() { + + } + + private List parseClusterNodes(String nodes) { + List pairsList = new LinkedList<>(); + logger.info("elasticsearch cluster nodes: {}", nodes); + String[] nodesSplit = nodes.split(","); + for (int i = 0; i < nodesSplit.length; i++) { + String node = nodesSplit[i]; + String host = node.split(":")[0]; + String port = node.split(":")[1]; + pairsList.add(new AddressPairs(host, Integer.valueOf(port))); + } + + return pairsList; + } + + class AddressPairs { + private String host; + private Integer port; + + public AddressPairs(String host, Integer port) { + this.host = host; + this.port = port; + } + } + + public boolean createIndex(String indexName, String indexType, Settings settings, XContentBuilder mappingBuilder) { + IndicesAdminClient adminClient = client.admin().indices(); + CreateIndexResponse response = adminClient.prepareCreate(indexName).setSettings(settings).addMapping(indexType, mappingBuilder).get(); + logger.info("create {} index with type of {} finished, isAcknowledged: {}", indexName, indexType, response.isAcknowledged()); + return response.isShardsAcked(); + } + + public boolean deleteIndex(String indexName) { + IndicesAdminClient adminClient = client.admin().indices(); + DeleteIndexResponse response = adminClient.prepareDelete(indexName).get(); + logger.info("delete {} index finished, isAcknowledged: {}", indexName, response.isAcknowledged()); + return response.isAcknowledged(); + } + + public boolean isExistsIndex(String indexName) { + IndicesAdminClient adminClient = client.admin().indices(); + IndicesExistsResponse response = adminClient.prepareExists(indexName).get(); + return response.isExists(); + } + + public SearchRequestBuilder prepareSearch(String indexName) { + return client.prepareSearch(indexName); + } + + public IndexRequestBuilder prepareIndex(String indexName, String id) { + return client.prepareIndex(indexName, "type", id); + } + + public UpdateRequestBuilder prepareUpdate(String indexName, String id) { + return client.prepareUpdate(indexName, "type", id); + } + + public GetRequestBuilder prepareGet(String indexName, String id) { + return client.prepareGet(indexName, "type", id); + } + + public MultiGetRequestBuilder prepareMultiGet() { + return client.prepareMultiGet(); + } + + public BulkRequestBuilder prepareBulk() { + return client.prepareBulk(); + } + + public void update(UpdateRequest updateRequest) { + try { + client.update(updateRequest).get(); + } catch (InterruptedException | ExecutionException e) { + logger.error(e.getMessage(), e); + } + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClientException.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClientException.java new file mode 100644 index 000000000..24eb74459 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClientException.java @@ -0,0 +1,34 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.elasticsearch; + +import org.skywalking.apm.collector.core.client.ClientException; + +/** + * @author peng-yongsheng + */ +public class ElasticSearchClientException extends ClientException { + public ElasticSearchClientException(String message) { + super(message); + } + + public ElasticSearchClientException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClient.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClient.java new file mode 100644 index 000000000..782148bbd --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClient.java @@ -0,0 +1,57 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.grpc; + +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import org.skywalking.apm.collector.core.client.Client; +import org.skywalking.apm.collector.core.client.ClientException; + +/** + * @author peng-yongsheng + */ +public class GRPCClient implements Client { + + private final String host; + + private final int port; + + private ManagedChannel channel; + + public GRPCClient(String host, int port) { + this.host = host; + this.port = port; + } + + @Override public void initialize() throws ClientException { + channel = ManagedChannelBuilder.forAddress(host, port).usePlaintext(true).build(); + } + + @Override public void shutdown() { + channel.shutdownNow(); + } + + public ManagedChannel getChannel() { + return channel; + } + + @Override public String toString() { + return host + ":" + port; + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClientException.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClientException.java new file mode 100644 index 000000000..a0322ea6a --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClientException.java @@ -0,0 +1,35 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.grpc; + +import org.skywalking.apm.collector.core.client.ClientException; + +/** + * @author peng-yongsheng + */ +public class GRPCClientException extends ClientException { + + public GRPCClientException(String message) { + super(message); + } + + public GRPCClientException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java new file mode 100644 index 000000000..74dd9b922 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java @@ -0,0 +1,119 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.h2; + +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import org.h2.util.IOUtils; +import org.skywalking.apm.collector.core.client.Client; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class H2Client implements Client { + + private final Logger logger = LoggerFactory.getLogger(H2Client.class); + + private Connection conn; + private String url; + private String userName; + private String password; + + public H2Client() { + this.url = "jdbc:h2:mem:collector"; + this.userName = ""; + this.password = ""; + } + + public H2Client(String url, String userName, String password) { + this.url = url; + this.userName = userName; + this.password = password; + } + + @Override public void initialize() throws H2ClientException { + try { + Class.forName("org.h2.Driver"); + conn = DriverManager. + getConnection(this.url, this.userName, this.password); + } catch (Exception e) { + throw new H2ClientException(e.getMessage(), e); + } + } + + @Override public void shutdown() { + IOUtils.closeSilently(conn); + } + + public Connection getConnection() throws H2ClientException { + return conn; + } + + public void execute(String sql) throws H2ClientException { + try (Statement statement = getConnection().createStatement()) { + statement.execute(sql); + statement.closeOnCompletion(); + } catch (SQLException e) { + throw new H2ClientException(e.getMessage(), e); + } + } + + public ResultSet executeQuery(String sql, Object[] params) throws H2ClientException { + logger.debug("execute query with result: {}", sql); + ResultSet rs; + PreparedStatement statement; + try { + statement = getConnection().prepareStatement(sql); + if (params != null) { + for (int i = 0; i < params.length; i++) { + statement.setObject(i + 1, params[i]); + } + } + rs = statement.executeQuery(); + statement.closeOnCompletion(); + } catch (SQLException e) { + throw new H2ClientException(e.getMessage(), e); + } + return rs; + } + + public boolean execute(String sql, Object[] params) throws H2ClientException { + logger.debug("execute insert/update/delete: {}", sql); + boolean flag; + Connection conn = getConnection(); + try (PreparedStatement statement = conn.prepareStatement(sql)) { + conn.setAutoCommit(true); + if (params != null) { + for (int i = 0; i < params.length; i++) { + statement.setObject(i + 1, params[i]); + } + } + flag = statement.execute(); + } catch (SQLException e) { + throw new H2ClientException(e.getMessage(), e); + } + return flag; + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/h2/H2ClientException.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/h2/H2ClientException.java new file mode 100644 index 000000000..45e4828a8 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/h2/H2ClientException.java @@ -0,0 +1,35 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.h2; + +import org.skywalking.apm.collector.core.client.ClientException; + +/** + * @author peng-yongsheng + */ +public class H2ClientException extends ClientException { + + public H2ClientException(String message) { + super(message); + } + + public H2ClientException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/redis/RedisClient.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/redis/RedisClient.java new file mode 100644 index 000000000..e75954743 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/redis/RedisClient.java @@ -0,0 +1,51 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.redis; + +import org.skywalking.apm.collector.core.client.Client; +import org.skywalking.apm.collector.core.client.ClientException; +import redis.clients.jedis.Jedis; + +/** + * @author peng-yongsheng + */ +public class RedisClient implements Client { + + private Jedis jedis; + + private final String host; + private final int port; + + public RedisClient(String host, int port) { + this.host = host; + this.port = port; + } + + @Override public void initialize() throws ClientException { + jedis = new Jedis(host, port); + } + + @Override public void shutdown() { + + } + + public void setex(String key, int seconds, String value) { + jedis.setex(key, seconds, value); + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/redis/RedisClientException.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/redis/RedisClientException.java new file mode 100644 index 000000000..0e3017225 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/redis/RedisClientException.java @@ -0,0 +1,35 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.redis; + +import org.skywalking.apm.collector.core.client.ClientException; + +/** + * @author peng-yongsheng + */ +public class RedisClientException extends ClientException { + + public RedisClientException(String message) { + super(message); + } + + public RedisClientException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClient.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClient.java new file mode 100644 index 000000000..a561e8b80 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClient.java @@ -0,0 +1,112 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.zookeeper; + +import java.io.IOException; +import java.util.List; +import org.apache.zookeeper.CreateMode; +import org.apache.zookeeper.KeeperException; +import org.apache.zookeeper.Watcher; +import org.apache.zookeeper.ZooKeeper; +import org.apache.zookeeper.data.ACL; +import org.apache.zookeeper.data.Stat; +import org.skywalking.apm.collector.core.client.Client; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class ZookeeperClient implements Client { + + private final Logger logger = LoggerFactory.getLogger(ZookeeperClient.class); + + private ZooKeeper zk; + + private final String hostPort; + private final int sessionTimeout; + private final Watcher watcher; + + public ZookeeperClient(String hostPort, int sessionTimeout, Watcher watcher) { + this.hostPort = hostPort; + this.sessionTimeout = sessionTimeout; + this.watcher = watcher; + } + + @Override public void initialize() throws ZookeeperClientException { + try { + zk = new ZooKeeper(hostPort, sessionTimeout, watcher); + } catch (IOException e) { + throw new ZookeeperClientException(e.getMessage(), e); + } + } + + @Override public void shutdown() { + + } + + public void create(final String path, byte data[], List acl, + CreateMode createMode) throws ZookeeperClientException { + try { + zk.create(path, data, acl, createMode); + } catch (KeeperException | InterruptedException e) { + throw new ZookeeperClientException(e.getMessage(), e); + } + } + + public Stat exists(final String path, boolean watch) throws ZookeeperClientException { + try { + return zk.exists(path, watch); + } catch (KeeperException | InterruptedException e) { + throw new ZookeeperClientException(e.getMessage(), e); + } + } + + public void delete(final String path, int version) throws ZookeeperClientException { + try { + zk.delete(path, version); + } catch (KeeperException | InterruptedException e) { + throw new ZookeeperClientException(e.getMessage(), e); + } + } + + public byte[] getData(String path, boolean watch, Stat stat) throws ZookeeperClientException { + try { + return zk.getData(path, watch, stat); + } catch (KeeperException | InterruptedException e) { + throw new ZookeeperClientException(e.getMessage(), e); + } + } + + public Stat setData(final String path, byte data[], int version) throws ZookeeperClientException { + try { + return zk.setData(path, data, version); + } catch (KeeperException | InterruptedException e) { + throw new ZookeeperClientException(e.getMessage(), e); + } + } + + public List getChildren(final String path, boolean watch) throws ZookeeperClientException { + try { + return zk.getChildren(path, watch); + } catch (KeeperException | InterruptedException e) { + throw new ZookeeperClientException(e.getMessage(), e); + } + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClientException.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClientException.java new file mode 100644 index 000000000..cdf80f33f --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClientException.java @@ -0,0 +1,34 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.zookeeper; + +import org.skywalking.apm.collector.core.client.ClientException; + +/** + * @author peng-yongsheng + */ +public class ZookeeperClientException extends ClientException { + public ZookeeperClientException(String message) { + super(message); + } + + public ZookeeperClientException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/util/PathUtils.java b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/util/PathUtils.java new file mode 100644 index 000000000..ccf8d37e4 --- /dev/null +++ b/apm-collector/apm-collector-component/client-component/src/main/java/org/skywalking/apm/collector/client/zookeeper/util/PathUtils.java @@ -0,0 +1,34 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.client.zookeeper.util; + +/** + * @author peng-yongsheng + */ +public class PathUtils { + + public static String convertKey2Path(String key) { + String[] keys = key.split("\\."); + StringBuilder pathBuilder = new StringBuilder(); + for (String subPath : keys) { + pathBuilder.append("/").append(subPath); + } + return pathBuilder.toString(); + } +} diff --git a/apm-collector/apm-collector-component/pom.xml b/apm-collector/apm-collector-component/pom.xml index 353200d10..b7dbe6e0a 100644 --- a/apm-collector/apm-collector-component/pom.xml +++ b/apm-collector/apm-collector-component/pom.xml @@ -19,4 +19,12 @@ stream-component remote-component + + + + org.skywalking + apm-collector-core + ${project.version} + + \ No newline at end of file diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/CollectorException.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/CollectorException.java new file mode 100644 index 000000000..ba459eeeb --- /dev/null +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/CollectorException.java @@ -0,0 +1,33 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.core; + +/** + * @author peng-yongsheng + */ +public class CollectorException extends Exception { + + public CollectorException(String message) { + super(message); + } + + public CollectorException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/Client.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/Client.java new file mode 100644 index 000000000..d6f4141eb --- /dev/null +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/Client.java @@ -0,0 +1,28 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.core.client; + +/** + * @author peng-yongsheng + */ +public interface Client { + void initialize() throws ClientException; + + void shutdown(); +} diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/ClientException.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/ClientException.java new file mode 100644 index 000000000..b3088011c --- /dev/null +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/ClientException.java @@ -0,0 +1,34 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.core.client; + +import org.skywalking.apm.collector.core.CollectorException; + +/** + * @author peng-yongsheng + */ +public abstract class ClientException extends CollectorException { + public ClientException(String message) { + super(message); + } + + public ClientException(String message, Throwable cause) { + super(message, cause); + } +}