Removed connectable interface

This commit is contained in:
Walter Dal Mut
2014-09-13 00:10:45 +02:00
parent 723ff0ecd2
commit 1a06353b24
8 changed files with 256 additions and 48 deletions
+1
View File
@@ -10,6 +10,7 @@ before_install:
- wget http://s3.amazonaws.com/influxdb/influxdb_latest_amd64.deb
- sudo useradd influxdb
- sudo dpkg -i influxdb_latest_amd64.deb
- sudo cp ./scripts/influxdb_conf.toml /opt/influxdb/shared/config.toml
- travis_retry sudo service influxdb restart
- sudo service influxdb status
+187
View File
@@ -0,0 +1,187 @@
# Welcome to the InfluxDB configuration file.
# If hostname (on the OS) doesn't return a name that can be resolved by the other
# systems in the cluster, you'll have to set the hostname to an IP or something
# that can be resolved here.
# hostname = ""
bind-address = "0.0.0.0"
# Once every 24 hours InfluxDB will report anonymous data to m.influxdb.com
# The data includes raft name (random 8 bytes), os, arch and version
# We don't track ip addresses of servers reporting. This is only used
# to track the number of instances running and the versions which
# is very helpful for us.
# Change this option to true to disable reporting.
reporting-disabled = false
[logging]
# logging level can be one of "debug", "info", "warn" or "error"
level = "info"
file = "/opt/influxdb/shared/log.txt" # stdout to log to standard out
# Configure the admin server
[admin]
port = 8083 # binding is disabled if the port isn't set
assets = "/opt/influxdb/current/admin"
# Configure the http api
[api]
port = 8086 # binding is disabled if the port isn't set
# ssl-port = 8084 # Ssl support is enabled if you set a port and cert
# ssl-cert = /path/to/cert.pem
# connections will timeout after this amount of time. Ensures that clients that misbehave
# and keep alive connections they don't use won't end up connection a million times.
# However, if a request is taking longer than this to complete, could be a problem.
read-timeout = "5s"
[input_plugins]
# Configure the graphite api
[input_plugins.graphite]
enabled = false
# port = 2003
# database = "" # store graphite data in this database
# udp_enabled = true # enable udp interface on the same port as the tcp interface
# Configure the udp api
#[input_plugins.udp]
#enabled = true
#port = 4444
# database = ""
# Configure multiple udp apis each can write to separate db. Just
# repeat the following section to enable multiple udp apis on
# different ports.
[[input_plugins.udp_servers]] # array of tables
enabled = true
port = 5551
database = "mine"
# Raft configuration
[raft]
# The raft port should be open between all servers in a cluster.
# However, this port shouldn't be accessible from the internet.
port = 8090
# Where the raft logs are stored. The user running InfluxDB will need read/write access.
dir = "/opt/influxdb/shared/data/raft"
# election-timeout = "1s"
[storage]
dir = "/opt/influxdb/shared/data/db"
# How many requests to potentially buffer in memory. If the buffer gets filled then writes
# will still be logged and once the local storage has caught up (or compacted) the writes
# will be replayed from the WAL
write-buffer-size = 10000
# the engine to use for new shards, old shards will continue to use the same engine
default-engine = "rocksdb"
# The default setting on this is 0, which means unlimited. Set this to something if you want to
# limit the max number of open files. max-open-files is per shard so this * that will be max.
max-open-shards = 0
# The default setting is 100. This option tells how many points will be fetched from LevelDb before
# they get flushed into backend.
point-batch-size = 100
# The number of points to batch in memory before writing them to leveldb. Lowering this number will
# reduce the memory usage, but will result in slower writes.
write-batch-size = 5000000
# The server will check this often for shards that have expired that should be cleared.
retention-sweep-period = "10m"
[storage.engines.leveldb]
# Maximum mmap open files, this will affect the virtual memory used by
# the process
max-open-files = 1000
# LRU cache size, LRU is used by leveldb to store contents of the
# uncompressed sstables. You can use `m` or `g` prefix for megabytes
# and gigabytes, respectively.
lru-cache-size = "200m"
[storage.engines.rocksdb]
# Maximum mmap open files, this will affect the virtual memory used by
# the process
max-open-files = 1000
# LRU cache size, LRU is used by rocksdb to store contents of the
# uncompressed sstables. You can use `m` or `g` prefix for megabytes
# and gigabytes, respectively.
lru-cache-size = "200m"
[storage.engines.hyperleveldb]
# Maximum mmap open files, this will affect the virtual memory used by
# the process
max-open-files = 1000
# LRU cache size, LRU is used by rocksdb to store contents of the
# uncompressed sstables. You can use `m` or `g` prefix for megabytes
# and gigabytes, respectively.
lru-cache-size = "200m"
[storage.engines.lmdb]
map-size = "100g"
[cluster]
# A comma separated list of servers to seed
# this server. this is only relevant when the
# server is joining a new cluster. Otherwise
# the server will use the list of known servers
# prior to shutting down. Any server can be pointed to
# as a seed. It will find the Raft leader automatically.
# Here's an example. Note that the port on the host is the same as the raft port.
# seed-servers = ["hosta:8090","hostb:8090"]
# Replication happens over a TCP connection with a Protobuf protocol.
# This port should be reachable between all servers in a cluster.
# However, this port shouldn't be accessible from the internet.
protobuf_port = 8099
protobuf_timeout = "2s" # the write timeout on the protobuf conn any duration parseable by time.ParseDuration
protobuf_heartbeat = "200ms" # the heartbeat interval between the servers. must be parseable by time.ParseDuration
protobuf_min_backoff = "1s" # the minimum backoff after a failed heartbeat attempt
protobuf_max_backoff = "10s" # the maxmimum backoff after a failed heartbeat attempt
# How many write requests to potentially buffer in memory per server. If the buffer gets filled then writes
# will still be logged and once the server has caught up (or come back online) the writes
# will be replayed from the WAL
write-buffer-size = 1000
# the maximum number of responses to buffer from remote nodes, if the
# expected number of responses exceed this number then querying will
# happen sequentially and the buffer size will be limited to this
# number
max-response-buffer-size = 100
# When queries get distributed out to shards, they go in parallel. This means that results can get buffered
# in memory since results will come in any order, but have to be processed in the correct time order.
# Setting this higher will give better performance, but you'll need more memory. Setting this to 1 will ensure
# that you don't need to buffer in memory, but you won't get the best performance.
concurrent-shard-query-limit = 10
[wal]
dir = "/opt/influxdb/shared/data/wal"
flush-after = 1000 # the number of writes after which wal will be flushed, 0 for flushing on every write
bookmark-after = 1000 # the number of writes after which a bookmark will be created
# the number of writes after which an index entry is created pointing
# to the offset of the first request, default to 1k
index-after = 1000
# the number of requests per one log file, if new requests came in a
# new log file will be created
requests-per-logfile = 10000
@@ -1,8 +0,0 @@
<?php
namespace InfluxDB\Adapter;
interface ConnectableInterface
{
public function connect();
public function disconnect();
}
+5
View File
@@ -42,6 +42,11 @@ class GuzzleAdapter implements AdapterInterface, QueryableInterface
}
$endpoint = $this->options->getHttpSeriesEndpoint();
try {
return $this->httpClient->get($endpoint, $options)->json();
} catch (\Exception $e) {
var_dump((string)$e->getResponse()->getBody(true));
die();
}
}
}
+4 -19
View File
@@ -3,35 +3,20 @@ namespace InfluxDB\Adapter;
use InfluxDB\Options;
class UdpAdapter implements AdapterInterface, ConnectableInterface
class UdpAdapter implements AdapterInterface
{
private $options;
private $socket;
public function __construct(Options $options)
{
$this->options = $options;
}
public function getSocket()
{
return $this->socket;
}
public function connect()
{
$this->socket = socket_create(AF_INET, SOCK_DGRAM, SOL_UDP);
return $this;
}
public function disconnect()
{
socket_close($this->getSocket());
}
public function send($message)
{
$socket = socket_create(AF_INET, SOCK_DGRAM, SOL_UDP);
$message = json_encode($message);
socket_sendto($this->getSocket(), $message, strlen($message), 0, $this->host, $this->port);
socket_sendto($socket, $message, strlen($message), 0, $this->options->getHost(), $this->options->getPort());
socket_close($socket);
}
}
-20
View File
@@ -20,26 +20,6 @@ class Client
return $this->adapter;
}
public function connect()
{
$result = false;
if ($this->getAdapter() instanceOf ConnectableInterface) {
$result = $this->getAdapter()->connect();
}
return $result;
}
public function disconnect()
{
$result = false;
if ($this->getAdapter() instanceOf ConnectableInterface) {
$result = $this->getAdapter()->disconnect();
}
return $result;
}
public function mark($name, array $values)
{
$data =[];
+51
View File
@@ -3,11 +3,13 @@ namespace InfluxDB;
use InfluxDB\Adapter\GuzzleAdapter as InfluxHttpAdapter;
use InfluxDB\Options;
use InfluxDB\Adapter\UdpAdapter;
use GuzzleHttp\Client as GuzzleHttpClient;
use crodas\InfluxPHP\Client as Crodas;
class ClientTest extends \PHPUnit_Framework_TestCase
{
private $rawOptions;
private $object;
private $options;
@@ -16,6 +18,7 @@ class ClientTest extends \PHPUnit_Framework_TestCase
public function setUp()
{
$options = include __DIR__ . '/../bootstrap.php';
$this->rawOptions = $options;
$client = new Crodas(
$options["tcp"]["host"],
@@ -30,6 +33,13 @@ class ClientTest extends \PHPUnit_Framework_TestCase
}
$client->createDatabase($options["tcp"]["database"]);
try {
$client->deleteDatabase($options["udp"]["database"]);
} catch (\Exception $e) {
// nothing...
}
$client->createDatabase($options["udp"]["database"]);
$this->anotherClient = $client;
$tcpOptions = $options["tcp"];
@@ -51,6 +61,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase
$this->object = $influx;
}
/**
* @group tcp
*/
public function testGuzzleHttpApiWorksCorrectly()
{
$this->object->mark("tcp.test", ["mark" => "element"]);
@@ -60,6 +73,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase
$this->assertEquals("element", $body[0]["points"][0][2]);
}
/**
* @group tcp
*/
public function testGuzzleHttpQueryApiWorksCorrectly()
{
$this->object->mark("tcp.test", ["mark" => "element"]);
@@ -71,6 +87,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase
$this->assertEquals("element", $body[0]["points"][0][2]);
}
/**
* @group tcp
*/
public function testGuzzleHttpQueryApiWithMultipleData()
{
$this->object->mark("tcp.test", ["mark" => "element"]);
@@ -83,6 +102,9 @@ class ClientTest extends \PHPUnit_Framework_TestCase
$this->assertEquals("tcp.test", $body[0]["name"]);
}
/**
* @group tcp
*/
public function testGuzzleHttpQueryApiWithTimePrecision()
{
$this->object->mark("tcp.test", ["mark" => "element"]);
@@ -92,4 +114,33 @@ class ClientTest extends \PHPUnit_Framework_TestCase
$this->assertCount(1, $body[0]["points"]);
$this->assertEquals("tcp.test", $body[0]["name"]);
}
/**
* @group udp
*/
public function testUdpIpWriteData()
{
$rawOptions = $this->rawOptions;
$options = new Options();
$options->setUsername($rawOptions["udp"]["username"]);
$options->setPassword($rawOptions["udp"]["password"]);
$options->setPort($rawOptions["udp"]["port"]);
$adapter = new UdpAdapter($options);
$object = new Client();
$object->setAdapter($adapter);
$object->mark("udp.test", ["mark" => "element"]);
$object->mark("udp.test", ["mark" => "element1"]);
$object->mark("udp.test", ["mark" => "element2"]);
$object->mark("udp.test", ["mark" => "element3"]);
// Wait UDP/IP message arrives
usleep(200e3);
$body = $this->object->query("select * from udp.test");
$this->assertCount(4, $body[0]["points"]);
$this->assertEquals("udp.test", $body[0]["name"]);
}
}
+8 -1
View File
@@ -7,5 +7,12 @@ return [
"database" => "mine",
"username" => "root",
"password" => "root",
]
],
"udp" => [
"host" => "localhost",
"port" => 5551,
"database" => "mine",
"username" => "root",
"password" => "root"
],
];