77import maya
88import requests
99from influxdb import InfluxDBClient
10+ from influxdb_client import InfluxDBClient , Point
11+ from influxdb_client .client .write_api import SYNCHRONOUS
1012
1113
1214def retrieve_paginated_data (
@@ -31,7 +33,7 @@ def retrieve_paginated_data(
3133 return results
3234
3335
34- def store_series (connection , series , metrics , rate_data ):
36+ def store_series (connection , version , org , bucket , series , metrics , rate_data ):
3537
3638 agile_data = rate_data .get ('agile_unit_rates' , [])
3739 agile_rates = {
@@ -131,8 +133,10 @@ def tags_for_measurement(measurement):
131133 }
132134 for measurement in metrics
133135 ]
134- connection .write_points (measurements )
135-
136+ if (version == 2 ):
137+ connection .write (bucket , org , measurements )
138+ elif (version == 1 ):
139+ connection .write_points (measurements )
136140
137141@click .command ()
138142@click .option (
@@ -147,13 +151,28 @@ def cmd(config_file, from_date, to_date):
147151 config = ConfigParser ()
148152 config .read (config_file )
149153
150- influx = InfluxDBClient (
151- host = config .get ('influxdb' , 'host' , fallback = "localhost" ),
152- port = config .getint ('influxdb' , 'port' , fallback = 8086 ),
153- username = config .get ('influxdb' , 'user' , fallback = "" ),
154- password = config .get ('influxdb' , 'password' , fallback = "" ),
155- database = config .get ('influxdb' , 'database' , fallback = "energy" ),
156- )
154+ org = config .get ('influxdb' , 'org' , fallback = "" )
155+ database = config .get ('influxdb' , 'database' , fallback = "energy" )
156+ influx_version = (config .getint ('influxdb' , 'version' , fallback = 2 ))
157+
158+ if (influx_version == 2 ):
159+ influx = InfluxDBClient (
160+ url = config .get ('influxdb' , 'url' , fallback = "http://localhost:8086" ),
161+ token = config .get ('influxdb' , 'token' , fallback = "" ),
162+ org = org ,
163+ )
164+ write_api = influx .write_api (write_options = SYNCHRONOUS )
165+ elif (influx_version == 1 ):
166+ influx = InfluxDBClient (
167+ host = config .get ('influxdb' , 'host' , fallback = "localhost" ),
168+ port = config .getint ('influxdb' , 'port' , fallback = 8086 ),
169+ username = config .get ('influxdb' , 'user' , fallback = "" ),
170+ password = config .get ('influxdb' , 'password' , fallback = "" ),
171+ database = database ,
172+ )
173+ write_api = influx
174+ else :
175+ raise click .ClickException ('Influx version not supported' )
157176
158177 api_key = config .get ('octopus' , 'api_key' )
159178 if not api_key :
@@ -231,7 +250,7 @@ def cmd(config_file, from_date, to_date):
231250 api_key , agile_url , from_iso , to_iso
232251 )
233252 click .echo (f' { len (rate_data ["electricity" ]["agile_unit_rates" ])} rates.' )
234- store_series (influx , 'electricity' , e_consumption , rate_data ['electricity' ])
253+ store_series (write_api , influx_version , org , database , 'electricity' , e_consumption , rate_data ['electricity' ])
235254
236255 click .echo (
237256 f'Retrieving gas data for { from_iso } until { to_iso } ...' ,
@@ -241,8 +260,8 @@ def cmd(config_file, from_date, to_date):
241260 api_key , g_url , from_iso , to_iso
242261 )
243262 click .echo (f' { len (g_consumption )} readings.' )
244- store_series (influx , 'gas' , g_consumption , rate_data ['gas' ])
263+ store_series (write_api , influx_version , org , database , 'gas' , g_consumption , rate_data ['gas' ])
245264
246265
247266if __name__ == '__main__' :
248- cmd ()
267+ cmd ()
0 commit comments