Skip to content
This repository has been archived by the owner on Feb 26, 2024. It is now read-only.

Commit

Permalink
Update for Influx v2
Browse files Browse the repository at this point in the history
  • Loading branch information
DragonQ committed Mar 9, 2021
1 parent 7ec20e8 commit 814160b
Show file tree
Hide file tree
Showing 4 changed files with 1,782 additions and 17 deletions.
45 changes: 32 additions & 13 deletions app/octopus_to_influxdb.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@
import maya
import requests
from influxdb import InfluxDBClient
from influxdb_client import InfluxDBClient, Point
from influxdb_client.client.write_api import SYNCHRONOUS


def retrieve_paginated_data(
Expand All @@ -31,7 +33,7 @@ def retrieve_paginated_data(
return results


def store_series(connection, series, metrics, rate_data):
def store_series(connection, version, org, bucket, series, metrics, rate_data):

agile_data = rate_data.get('agile_unit_rates', [])
agile_rates = {
Expand Down Expand Up @@ -114,8 +116,10 @@ def tags_for_measurement(measurement):
}
for measurement in metrics
]
connection.write_points(measurements)

if (version == 2):
connection.write(bucket, org, measurements)
elif (version == 1):
connection.write_points(measurements)

@click.command()
@click.option(
Expand All @@ -130,13 +134,28 @@ def cmd(config_file, from_date, to_date):
config = ConfigParser()
config.read(config_file)

influx = InfluxDBClient(
host=config.get('influxdb', 'host', fallback="localhost"),
port=config.getint('influxdb', 'port', fallback=8086),
username=config.get('influxdb', 'user', fallback=""),
password=config.get('influxdb', 'password', fallback=""),
database=config.get('influxdb', 'database', fallback="energy"),
)
org=config.get('influxdb', 'org', fallback="")
database=config.get('influxdb', 'database', fallback="energy")
influx_version=(config.getint('influxdb', 'version', fallback=2))

if (influx_version == 2):
influx = InfluxDBClient(
url=config.get('influxdb', 'url', fallback="http://localhost:8086"),
token=config.get('influxdb', 'token', fallback=""),
org=org,
)
write_api = influx.write_api(write_options=SYNCHRONOUS)
elif (influx_version == 1):
influx = InfluxDBClient(
host=config.get('influxdb', 'host', fallback="localhost"),
port=config.getint('influxdb', 'port', fallback=8086),
username=config.get('influxdb', 'user', fallback=""),
password=config.get('influxdb', 'password', fallback=""),
database=database,
)
write_api = influx
else:
raise click.ClickException('Influx version not supported')

api_key = config.get('octopus', 'api_key')
if not api_key:
Expand Down Expand Up @@ -214,7 +233,7 @@ def cmd(config_file, from_date, to_date):
api_key, agile_url, from_iso, to_iso
)
click.echo(f' {len(rate_data["electricity"]["agile_unit_rates"])} rates.')
store_series(influx, 'electricity', e_consumption, rate_data['electricity'])
store_series(write_api, influx_version, org, database, 'electricity', e_consumption, rate_data['electricity'])

click.echo(
f'Retrieving gas data for {from_iso} until {to_iso}...',
Expand All @@ -224,8 +243,8 @@ def cmd(config_file, from_date, to_date):
api_key, g_url, from_iso, to_iso
)
click.echo(f' {len(g_consumption)} readings.')
store_series(influx, 'gas', g_consumption, rate_data['gas'])
store_series(write_api, influx_version, org, database, 'gas', g_consumption, rate_data['gas'])


if __name__ == '__main__':
cmd()
cmd()
Loading

0 comments on commit 814160b

Please sign in to comment.