InfluxDB.Client
5.1.0
See the version list below for details.
dotnet add package InfluxDB.Client --version 5.1.0
NuGet\Install-Package InfluxDB.Client -Version 5.1.0
<PackageReference Include="InfluxDB.Client" Version="5.1.0" />
<PackageVersion Include="InfluxDB.Client" Version="5.1.0" />
<PackageReference Include="InfluxDB.Client" />
paket add InfluxDB.Client --version 5.1.0
#r "nuget: InfluxDB.Client, 5.1.0"
#:package InfluxDB.Client@5.1.0
#addin nuget:?package=InfluxDB.Client&version=5.1.0
#tool nuget:?package=InfluxDB.Client&version=5.1.0
InfluxDB.Client
The reference client that allows query, write and management (bucket, organization, users) for the InfluxDB 2.x.
Documentation
This section contains links to the client library documentation.
Features
- Querying data using Flux language
- Writing data using
- Delete data
- InfluxDB 2.x Management API
- sources, buckets
- tasks
- authorizations
- health check
- Advanced Usage
Queries
For querying data we use QueryApi that allow perform asynchronous, streaming, synchronous and also use raw query response.
Asynchronous Query
The asynchronous query is not intended for large query results because the Flux response can be potentially unbound.
using System;
using System.Threading.Tasks;
using InfluxDB.Client;
namespace Examples
{
public static class AsynchronousQuery
{
private static readonly string Token = "";
public static async Task Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
var flux = "from(bucket:\"temperature-sensors\") |> range(start: 0)";
var queryApi = client.GetQueryApi();
//
// QueryData
//
var tables = await queryApi.QueryAsync(flux, "org_id");
tables.ForEach(table =>
{
table.Records.ForEach(record =>
{
Console.WriteLine($"{record.GetTime()}: {record.GetValueByKey("_value")}");
});
});
}
}
}
The asynchronous query offers a possibility map FluxRecords to POCO:
using System;
using System.Threading.Tasks;
using InfluxDB.Client;
using InfluxDB.Client.Core;
namespace Examples
{
public static class AsynchronousQuery
{
private static readonly string Token = "";
public static async Task Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
var flux = "from(bucket:\"temperature-sensors\") |> range(start: 0)";
var queryApi = client.GetQueryApi();
//
// QueryData
//
var temperatures = await queryApi.QueryAsync<Temperature>(flux, "org_id");
temperatures.ForEach(temperature =>
{
Console.WriteLine($"{temperature.Location}: {temperature.Value} at {temperature.Time}");
});
}
[Measurement("temperature")]
private class Temperature
{
[Column("location", IsTag = true)] public string Location { get; set; }
[Column("value")] public double Value { get; set; }
[Column(IsTimestamp = true)] public DateTime Time { get; set; }
}
}
}
Streaming Query
The Streaming query offers possibility to process unbound query and allow user to handle exceptions, stop receiving more results and notify that all data arrived.
using System;
using System.Threading.Tasks;
using InfluxDB.Client;
namespace Examples
{
public static class StreamingQuery
{
private static readonly string Token = "";
public static async Task Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
var flux = "from(bucket:\"temperature-sensors\") |> range(start: 0)";
var queryApi = client.GetQueryApi();
//
// QueryData
//
await queryApi.QueryAsync(flux, record =>
{
//
// The callback to consume a FluxRecord.
//
Console.WriteLine($"{record.GetTime()}: {record.GetValueByKey("_value")}");
}, exception =>
{
//
// The callback to consume any error notification.
//
Console.WriteLine($"Error occurred: {exception.Message}");
}, () =>
{
//
// The callback to consume a notification about successfully end of stream.
//
Console.WriteLine("Query completed");
}, "org_id");
}
}
}
And there is also a possibility map FluxRecords to POCO:
using System;
using System.Threading.Tasks;
using InfluxDB.Client;
using InfluxDB.Client.Core;
namespace Examples
{
public static class StreamingQuery
{
private static readonly string Token = "";
public static async Task Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
var flux = "from(bucket:\"temperature-sensors\") |> range(start: 0)";
var queryApi = client.GetQueryApi();
//
// QueryData
//
await queryApi.QueryAsync<Temperature>(flux, temperature =>
{
//
// The callback to consume a FluxRecord mapped to POCO.
//
Console.WriteLine($"{temperature.Location}: {temperature.Value} at {temperature.Time}");
}, org: "org_id");
}
[Measurement("temperature")]
private class Temperature
{
[Column("location", IsTag = true)] public string Location { get; set; }
[Column("value")] public double Value { get; set; }
[Column(IsTimestamp = true)] public DateTime Time { get; set; }
}
}
}
Raw Query
The Raw query allows direct processing original CSV response:
using System;
using System.Threading.Tasks;
using InfluxDB.Client;
namespace Examples
{
public static class RawQuery
{
private static readonly string Token = "";
public static async Task Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
var flux = "from(bucket:\"temperature-sensors\") |> range(start: 0)";
var queryApi = client.GetQueryApi();
//
// QueryData
//
var csv = await queryApi.QueryRawAsync(flux, org: "org_id");
Console.WriteLine($"CSV response: {csv}");
}
}
}
The Streaming version allows processing line by line:
using System;
using System.Threading.Tasks;
using InfluxDB.Client;
namespace Examples
{
public static class RawQueryAsynchronous
{
private static readonly string Token = "";
public static async Task Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
var flux = "from(bucket:\"temperature-sensors\") |> range(start: 0)";
var queryApi = client.GetQueryApi();
//
// QueryData
//
await queryApi.QueryRawAsync(flux, line =>
{
//
// The callback to consume a line of CSV response
//
Console.WriteLine($"Response: {line}");
}, org: "org_id");
}
}
}
Synchronous query
The synchronous query is not intended for large query results because the response can be potentially unbound.
using System;
using InfluxDB.Client;
namespace Examples
{
public static class SynchronousQuery
{
public static void Main()
{
using var client = new InfluxDBClient("http://localhost:9999", "my-token");
const string query = "from(bucket:\"my-bucket\") |> range(start: 0)";
//
// QueryData
//
var queryApi = client.GetQueryApiSync();
var tables = queryApi.QuerySync(query, "my-org");
//
// Process results
//
tables.ForEach(table =>
{
table.Records.ForEach(record =>
{
Console.WriteLine($"{record.GetTime()}: {record.GetValueByKey("_value")}");
});
});
}
}
}
Writes
For writing data we use WriteApi or WriteApiAsync which is simplified version of WriteApi without batching support.
WriteApi supports:
- writing data using InfluxDB Line Protocol, Data Point, POCO
- use batching for writes
- produces events that allow user to be notified and react to this events
WriteSuccessEvent- published when arrived the success response from serverWriteErrorEvent- published when occurs a unhandled exception from serverWriteRetriableErrorEvent- published when occurs a retriable error from serverWriteRuntimeExceptionEvent- published when occurs a runtime exception in background batch processing
- use GZIP compression for data
The writes are processed in batches which are configurable by WriteOptions:
| Property | Description | Default Value |
|---|---|---|
| BatchSize | the number of data point to collect in batch | 1000 |
| FlushInterval | the number of milliseconds before the batch is written | 1000 |
| JitterInterval | the number of milliseconds to increase the batch flush interval by a random amount | 0 |
| RetryInterval | the number of milliseconds to retry unsuccessful write. The retry interval is used when the InfluxDB server does not specify "Retry-After" header. | 5000 |
| MaxRetries | the number of max retries when write fails | 3 |
| MaxRetryDelay | the maximum delay between each retry attempt in milliseconds | 125_000 |
| ExponentialBase | the base for the exponential retry delay, the next delay is computed using random exponential backoff as a random value within the interval retryInterval * exponentialBase^(attempts-1) and retryInterval * exponentialBase^(attempts). Example for retryInterval=5_000, exponentialBase=2, maxRetryDelay=125_000, maxRetries=5 Retry delays are random distributed values within the ranges of [5_000-10_000, 10_000-20_000, 20_000-40_000, 40_000-80_000, 80_000-125_000] |
2 |
Writing data
By POCO
Write Measurement into specified bucket:
using System;
using InfluxDB.Client;
using InfluxDB.Client.Api.Domain;
using InfluxDB.Client.Core;
namespace Examples
{
public static class WritePoco
{
private static readonly string Token = "";
public static void Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
//
// Write Data
//
using (var writeApi = client.GetWriteApi())
{
//
// Write by POCO
//
var temperature = new Temperature {Location = "south", Value = 62D, Time = DateTime.UtcNow};
writeApi.WriteMeasurement(temperature, WritePrecision.Ns, "bucket_name", "org_id");
}
}
[Measurement("temperature")]
private class Temperature
{
[Column("location", IsTag = true)] public string Location { get; set; }
[Column("value")] public double Value { get; set; }
[Column(IsTimestamp = true)] public DateTime Time { get; set; }
}
}
}
By Data Point
Write Data point into specified bucket:
using System;
using InfluxDB.Client;
using InfluxDB.Client.Api.Domain;
using InfluxDB.Client.Writes;
namespace Examples
{
public static class WriteDataPoint
{
private static readonly string Token = "";
public static void Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
//
// Write Data
//
using (var writeApi = client.GetWriteApi())
{
//
// Write by Data Point
var point = PointData.Measurement("temperature")
.Tag("location", "west")
.Field("value", 55D)
.Timestamp(DateTime.UtcNow.AddSeconds(-10), WritePrecision.Ns);
writeApi.WritePoint(point, "bucket_name", "org_id");
}
}
}
}
DataPoint Builder Immutability: The builder is immutable therefore won't have side effect when using for building multiple point with single builder.
using System;
using InfluxDB.Client;
using InfluxDB.Client.Api.Domain;
using InfluxDB.Client.Writes;
namespace Examples
{
public static class WriteDataPoint
{
private static readonly string Token = "";
public static void Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
//
// Write Data
//
using (var writeApi = client.GetWriteApi())
{
//
// Write by Data Point
var builder = PointData.Measurement("temperature")
.Tag("location", "west");
var pointA = builder
.Field("value", 55D)
.Timestamp(DateTime.UtcNow.AddSeconds(-10), WritePrecision.Ns);
writeApi.WritePoint(pointA, "bucket_name", "org_id");
var pointB = builder
.Field("age", 32)
.Timestamp(DateTime.UtcNow, WritePrecision.Ns);
writeApi.WritePoint(pointB, "bucket_name", "org_id");
}
}
}
}
By LineProtocol
Write Line Protocol record into specified bucket:
using InfluxDB.Client;
using InfluxDB.Client.Api.Domain;
namespace Examples
{
public static class WriteLineProtocol
{
private static readonly string Token = "";
public static void Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
//
// Write Data
//
using (var writeApi = client.GetWriteApi())
{
//
//
// Write by LineProtocol
//
writeApi.WriteRecord("temperature,location=north value=60.0", WritePrecision.Ns,"bucket_name", "org_id");
}
}
}
}
Using WriteApiAsync
using System;
using System.Threading.Tasks;
using InfluxDB.Client;
using InfluxDB.Client.Api.Domain;
using InfluxDB.Client.Core;
using InfluxDB.Client.Writes;
namespace Examples
{
public static class WriteApiAsyncExample
{
[Measurement("temperature")]
private class Temperature
{
[Column("location", IsTag = true)] public string Location { get; set; }
[Column("value")] public double Value { get; set; }
[Column(IsTimestamp = true)] public DateTime Time { get; set; }
}
public static async Task Main()
{
using var client = new InfluxDBClient("http://localhost:8086",
"my-user", "my-password");
//
// Write Data
//
var writeApiAsync = client.GetWriteApiAsync();
//
//
// Write by LineProtocol
//
await writeApiAsync.WriteRecordAsync("temperature,location=north value=60.0", WritePrecision.Ns,
"my-bucket", "my-org");
//
//
// Write by Data Point
//
var point = PointData.Measurement("temperature")
.Tag("location", "west")
.Field("value", 55D)
.Timestamp(DateTime.UtcNow.AddSeconds(-10), WritePrecision.Ns);
await writeApiAsync.WritePointAsync(point, "my-bucket", "my-org");
//
// Write by POCO
//
var temperature = new Temperature {Location = "south", Value = 62D, Time = DateTime.UtcNow};
await writeApiAsync.WriteMeasurementAsync(temperature, WritePrecision.Ns, "my-bucket", "my-org");
//
// Check written data
//
var tables = await influxDbClient.GetQueryApi()
.QueryAsync("from(bucket:\"my-bucket\") |> range(start: 0)", "my-org");
tables.ForEach(table =>
{
var fluxRecords = table.Records;
fluxRecords.ForEach(record =>
{
Console.WriteLine($"{record.GetTime()}: {record.GetValue()}");
});
});
}
}
}
Default Tags
Sometimes is useful to store same information in every measurement e.g. hostname, location, customer.
The client is able to use static value, app settings or env variable as a tag value.
The expressions:
California Miner- static value${version}- application settings${env.hostname}- environment property
Via Configuration file
In a configuration file you are able to specify default tags by tags element.
<?xml version="1.0" encoding="utf-8"?>
<configuration>
<configSections>
<section name="influx2" type="InfluxDB.Client.Configurations.Influx2, InfluxDB.Client" />
</configSections>
<appSettings>
<add key="SensorVersion" value="v1.00"/>
</appSettings>
<influx2 url="http://localhost:8086"
org="my-org"
bucket="my-bucket"
token="my-token"
logLevel="BODY"
timeout="10s">
<tags>
<tag name="id" value="132-987-655"/>
<tag name="customer" value="California Miner"/>
<tag name="hostname" value="${env.Hostname}"/>
<tag name="sensor-version" value="${SensorVersion}"/>
</tags>
</influx2>
</configuration>
Via API
var options = new InfluxDBClientOptions(Url)
{
Token = token,
DefaultTags = new Dictionary<string, string>
{
{"id", "132-987-655"},
{"customer", "California Miner"},
}
};
options.AddDefaultTag("hostname", "${env.Hostname}")
options.AddDefaultTags(new Dictionary<string, string>{{ "sensor-version", "${SensorVersion}" }})
Both of configurations will produce the Line protocol:
mine-sensor,id=132-987-655,customer="California Miner",hostname=example.com,sensor-version=v1.00 altitude=10
Handle the Events
Events that can be handle by WriteAPI EventHandler are:
WriteSuccessEvent- for success response from serverWriteErrorEvent- for unhandled exception from serverWriteRetriableErrorEvent- for retriable error from serverWriteRuntimeExceptionEvent- for runtime exception in background batch processing
Number of events depends on number of data points to collect in batch. The batch size is configured by BatchSize option (default size is 1000) - in case
of one data point, event is handled for each point, independently on used writing method (even for mass writing of data like
WriteMeasurements, WritePoints and WriteRecords).
Events can be handled by register writeApi.EventHandler or by creating custom EventListener:
Register EventHandler
writeApi.EventHandler += (sender, eventArgs) =>
{
switch (eventArgs)
{
case WriteSuccessEvent successEvent:
string data = @event.LineProtocol;
//
// handle success response from server
// Console.WriteLine($"{data}");
//
break;
case WriteErrorEvent error:
string data = @error.LineProtocol;
string errorMessage = @error.Exception.Message;
//
// handle unhandled exception from server
//
// Console.WriteLine($"{data}");
// throw new Exception(errorMessage);
//
break;
case WriteRetriableErrorEvent error:
string data = @error.LineProtocol;
string errorMessage = @error.Exception.Message;
//
// handle retrievable error from server
//
// Console.WriteLine($"{data}");
// throw new Exception(errorMessage);
//
break;
case WriteRuntimeExceptionEvent error:
string errorMessage = @error.Exception.Message;
//
// handle runtime exception in background batch processing
// throw new Exception(errorMessage);
//
break;
}
};
//
// Write by LineProtocol
//
writeApi.WriteRecord("influxPoint,writeType=lineProtocol value=11.11" +
$" {DateTime.UtcNow.Subtract(EpochStart).Ticks * 100}", WritePrecision.Ns, "my-bucket", "my-org");
Custom EventListener
Advantage of using custom Event Listener is possibility of waiting on handled event between different writings - for more info see EventListener.
Delete Data
Delete data from specified bucket:
using InfluxDB.Client;
using InfluxDB.Client.Api.Domain;
namespace Examples
{
public static class WriteLineProtocol
{
private static readonly string Token = "";
public static void Main()
{
using var client = new InfluxDBClient("http://localhost:8086", Token);
//
// Delete data
//
await client.GetDeleteApi().Delete(DateTime.UtcNow.AddMinutes(-1), DateTime.Now, "", "bucket", "org");
}
}
}
Filter trace verbose
You can filter out verbose messages from InfluxDB.Client by using TraceListener.
using System;
using System.Diagnostics;
using InfluxDB.Client.Core;
namespace Examples
{
public static class MyProgram
{
public static void Main()
{
TraceListener ConsoleOutListener = new TextWriterTraceListener(Console.Out)
{
Filter = InfluxDBTraceFilter.SuppressInfluxVerbose(),
};
Trace.Listeners.Add(ConsoleOutListener);
// My code ...
}
}
}
Management API
The client has following management API:
| API endpoint | Description | Implementation |
|---|---|---|
| /api/v2/authorizations | Managing authorization data | AuthorizationsApi |
| /api/v2/buckets | Managing bucket data | BucketsApi |
| /api/v2/orgs | Managing organization data | OrganizationsApi |
| /api/v2/users | Managing user data | UsersApi |
| /api/v2/sources | Managing sources | SourcesApi |
| /api/v2/tasks | Managing one-off and recurring tasks | TasksApi |
| /api/v2/scrapers | Managing ScraperTarget data | ScraperTargetsApi |
| /api/v2/labels | Managing resource labels | LabelsApi |
| /api/v2/telegrafs | Managing telegraf config data | TelegrafsApi |
| /api/v2/setup | Managing onboarding setup | InfluxDBClient#OnBoarding() |