Files
energy-consumption-ingester/POC.ipynb
T
Brian Bjarke Jensen aa6ec36878
Python Code Quality / python-code-quality (push) Successful in 55s
Hourly Energy Data Sync / sync-energy-data (push) Successful in 17s
added full cli package and CI execution
2025-10-27 19:07:25 +01:00

894 lines
39 KiB
Plaintext

{
"cells": [
{
"cell_type": "markdown",
"id": "a2da9908",
"metadata": {},
"source": [
"# Setup"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "2e95933f",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Token request status: 200\n",
"Token response: {\"result\":\"eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJ0b2tlblR5cGUiOiJDdXN0b21lckFQSV9EYXRhQWNjZXNzIiwidG9rZW5pZCI6ImE2NzY1ODRmLWY1MzktNDI0Yy04MmZhLWM5MjNhMjNhOTVlYSIsIndlYkFwcCI6IkN1c3RvbWVyQXBwIiwidmVyc2lvbiI6IjIiLCJpZGVudGl0eVRva2VuIjoiWUpoaXlqMXRjMUZiTnBtWmNaM0tjdGUrbzNzQ1FxR3B1dWtEcVh4QzVveEhKMisvcTAycTdpWGlZZURTRm9jUjVmYVBZMHVCeGpiWnYrOE52SHUycnpFU2pFRDNUajBQbmM5QzBZY1FXb0hzTDRvVHliOWRTSysrVFR5NTdDWTR4bit1dTdYN2lYdTNyV1pLUTUzTkh5T2NRSjB3Y3ZtWURHUjVDT3lITW14WDZhMjFEUHdBOWFsRk53RlFXbGkxT0NuTnp3QUdQaXJZWnZLcTYxNW52aFBjczZaN0FWUlh5MnRqMTFPbWpuR1NQVFJYU2pvdjdpRkhZcmhTdTRiTFpDQXhDd0FHWHl4a0tPWENRRDZFUVdTUU92d0ZpMHdnSUZwVHBIQjJKaHlEVFRvM2FYU1RaSnh2MEpFZktvalNLWXUxSE54enpWcTE3NmVWei9UVW5JZTdUeTc0LzZNekFqZmplcWpFUTk3aXcvbnRwTFNzdGkwakZqdjB6V1hKRXRKMVNQYURmbG5ZdG9DcEtnZmdCVDc4RTNKcXVIZSszQjJLanUwL0xlUU4rMXpHU2cxeEhuUFZOdlcxanRabjhWbjhIaVVmNUhUK0tPZ3ZDNitqaWJQRWRwNXJvaWtSNkFjRzZHcS8xOThkL2d6bEo4MzJMVWVnR1BiTlBHVVBwZ2VQTmFtUlRSWXpnSk9GdmwzV0NGb1NXMUZhYU5UalBJL1F4Wjh4L05VVzlJZXFXbG1rb1VrMVhIS3lWSzBjdW1KWXlaelJBNHJKdGJ0QXVaaGdMUCtpUmtSUzlFMkcwOWhDQnhzOHJKWjRycWVvZmVITC9iQkQwYXIzd25tREFpUks3L1ovNmkvbXFJMEZHa3NUcXdYZzdBTTY0SldGQkNUcm1ZMk9hSlIwRHpSZmd3ekN5Y0NzRHRRUW1RNTZwNWVPRzJyRnlUUHh5YjEwS3p4a3FTbTZVNDExcEhlVllFbnR4cVloYUg4and2RWJxVFJBNEZ5VUZtQVVxQmhIVDFWNWpsRkFZL0U3bjUwejAxSE5WL21KNzVBcGNqektsaUx4YVBuQXdYUG43NmVVbHc4NUt4anJueEh5d3V6a00wVFhDSVRIZ0NrWmdYbkdxVEYwZE1TL3VhRE9Tbjl5NTY5SE9vb1FxM3VNZ0prV2RFaUdSMkUvcVErbFVNS05kVlNEL0F6Y0tLZkdhdkZUTko2WVhZOVozKzZKN2lpMWFDbGxoWFAwQ1VRRzhaOVpaMEMrU3JlN3FNSzVvaEpwY2FDOUdoZ1hVNklqZ0xZd2JTM0RGOUN3NGt3cmxHRVJSNURVZVNXREg2aFp5SVdEWWc2d1k4bHMxNGVlVjQ4bW5aTVVCSVE4cDZvNzhGSzFLUWc2V3BQZ3NqMGVOcXdBVW1HemNQcUs5UHJsdDc5VmZvK2pQSDlkUW1IaUJIbjFDcjROWWJpUEI4bllNN09JdFJpb3lHVmpyTi9lVjR1SGloc1c1ZzUxMG9pMjVxNzU0Y2VPbW10ck0ra0VzRGtoVEd1UzNxZml5ZUhqRDV0dk1leXU3eGs0c1gyeDdhVkdDeGV0eGlDalNmWkZ0R2JiWmY3VWhKWTAzNzBZbDYvZFRNeDJ4VFNlUzhqOUgwZUZJT1RBL3N3ZVNXUjdDMlc2MitWOC9SNERZT3lOdjJzWmRib2RIN3hKdk0yUnNxcHFRMUtqS1VsclpEUmRTMFk4ZDZTOGJuTTBGUC9tRXlDZkF1MmthRnE5d0hQcGZUdzJ2QVFIeTdMNnR5ZjNoSUNtOHQ1eGh4NEJwV1ZUUkVtQUVXZk5KQ3c2QXZlUFExVHY5TjVmQ1N6NXo5RlU2VmhxcUlQL1hWM2l3dUFQc3dVc2hpaG90ZTBPR2YrZlpyMEZHY04vSTJXYzZneDdIT3RYVitUWEdieGdvY2lianFWRTFyQTRweTJsYnpSQlo3dXVRc29OcW1YQlVkckF6U0V3MkFEZGg3cE1KWGpBRXF0cHZKOXluWUlqWWRpZXdocllEazRRNCtEekl5V01NRk5YOGVyakFER21PYnJvazB6aUtGd1Fwb3hCL2I2SU9GYzVwWjFvUm0zY0dNUjU2M3BPdmw0aTNIalhMbmxvbWkxQ29Hb1JPL203MGp2RXRMdXBvOGRRSGpFZXlEUzBmeGl1b0hHRWV1Uk12SGpjcVUranBadjloUjEwMXh1QVVGRjNaUFJ6NnhQb2JsZVR0N0xnSjZWMm9KZmJ4QmpPWEdLTEthdGNiUEkzMS84MWtSWlIyN0w5eVZ4NkRLays1ZjV4Zk5JdUxWd0pzeDNBRVBob2o4SERPVFRKSERMUlZpa1hOdGE4OTNZOWVaZ2xpRzdyMXZCaXRTZ2dwQTRlbDhjMHlIdGJDaG9Vb3N6SmsyMHZJWGphem1zRG5kTU5DVXluRGlWbEN4VDBROWRPMlc4Rlp5NERxRmxNN1BQVVZWZzRhMkRIUzhianZQR3dWZ0R6WDkrVzA3VEMrVVlpZ3k2Sk93YnkwMThqd2loQjlLb3d6M0Vqc1djTjF1RGRsYXZPbHNWbHFBU05tWUF0MmpnbVF1YnQzTUlrdFhObGNOMWl0MnA5NFJ3UE9qNmhjYm1OM3ZRaXNvVExyd1llK1RlciIsImh0dHA6Ly9zY2hlbWFzLnhtbHNvYXAub3JnL3dzLzIwMDUvMDUvaWRlbnRpdHkvY2xhaW1zL25hbWVpZGVudGlmaWVyIjoiTUlEOmUyZjE2NDE4LTlmZDItNGFkZi1hNDA2LWEwMDQ5ZTZjYzU4YSIsImh0dHA6Ly9zY2hlbWFzLnhtbHNvYXAub3JnL3dzLzIwMDUvMDUvaWRlbnRpdHkvY2xhaW1zL2dpdmVubmFtZSI6IkJyaWFuIEJqYXJrZSBKZW5zZW4iLCJsb2dpblR5cGUiOiJLZXlDYXJkIiwiYjNmIjoiRjRSQ1VHSjdoTUUyMlpUVzRnTGZha2crL0VXL2Y4MmZVVTFoNVFjMk1DWT0iLCJwaWQiOiJQSUQ6OTIwOC0yMDAyLTItNjAwNTA4Mzg2NzY4IiwidXNlcklkIjoiMTA0MzQ1OSIsImV4cCI6MTc2MTY0NzgxNywiaXNzIjoiRW5lcmdpbmV0IiwianRpIjoiYTY3NjU4NGYtZjUzOS00MjRjLTgyZmEtYzkyM2EyM2E5NWVhIiwidG9rZW5OYW1lIjoiZ2l0ZWEtYWNjZXNzLXRva2VuIiwiYXVkIjoiRW5lcmdpbmV0In0.zmb-Oy1Jn6u16dAodmdd3OApFbUVcTjetM5Z8x_evxQ\"}\n",
"Successfully obtained access token: eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJ0b2tlblR5c...\n"
]
}
],
"source": [
"import os\n",
"from dotenv import load_dotenv\n",
"import http.client\n",
"import json\n",
"from pydantic import AnyUrl\n",
"\n",
"load_dotenv()\n",
"ELOVERBLIK_API_TOKEN = os.getenv(\"ELOVERBLIK_API_TOKEN\")\n",
"\n",
"\n",
"server_url = AnyUrl(\"https://api.eloverblik.dk\")\n",
"\n",
"# Prepare JWT authentication\n",
"\n",
"conn = http.client.HTTPSConnection(server_url.host)\n",
"\n",
"headers = {\n",
" \"Authorization\": f\"Bearer {ELOVERBLIK_API_TOKEN}\",\n",
" \"Content-Type\": \"application/json\",\n",
" \"api-version\": \"1.0\",\n",
"}\n",
"\n",
"conn.request(\"GET\", \"/customerapi/api/token\", headers=headers)\n",
"\n",
"res = conn.getresponse()\n",
"print(f\"Token request status: {res.status}\")\n",
"\n",
"if res.status == 200:\n",
" token_response = res.read().decode(\"utf-8\")\n",
" print(f\"Token response: {token_response[:50]}...\")\n",
"\n",
" # Parse the response to get the access token\n",
" try:\n",
" token_data = json.loads(token_response)\n",
" access_token = token_data[\"result\"]\n",
" print(f\"Successfully obtained access token: {access_token[:50]}...\")\n",
" except Exception as e:\n",
" print(f\"Error parsing token response: {e}\")\n",
" access_token = None\n",
"else:\n",
" error_data = res.read().decode(\"utf-8\")\n",
" print(f\"Error getting token: {error_data}\")\n",
" access_token = None"
]
},
{
"cell_type": "code",
"execution_count": 33,
"id": "0a68b6bd",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Testing access token with metering points API...\n",
"Metering points API status: 200\n",
"Response: {\"result\":[{\"streetCode\":\"2724\",\"streetName\":\"Vejl...\n",
"Metering points API status: 200\n",
"Response: {\"result\":[{\"streetCode\":\"2724\",\"streetName\":\"Vejl...\n"
]
}
],
"source": [
"if not access_token:\n",
" print(\"No access token available to test with\")\n",
"else:\n",
" print(\"Testing access token with metering points API...\")\n",
"\n",
" # Create new connection for API calls\n",
" conn = http.client.HTTPSConnection(server_url.host)\n",
"\n",
" # Use the access token for API calls\n",
" api_headers = {\n",
" \"Authorization\": f\"Bearer {access_token}\",\n",
" \"Content-Type\": \"application/json\",\n",
" \"api-version\": \"1.0\",\n",
" }\n",
"\n",
" # Get metering points\n",
" from urllib.parse import urlencode\n",
"\n",
" params = {\n",
" \"includeAll\": True,\n",
" }\n",
"\n",
" conn.request(\n",
" \"GET\",\n",
" \"/customerapi/api/meteringpoints/meteringpoints?\" + urlencode(params),\n",
" headers=api_headers,\n",
" )\n",
"\n",
" res = conn.getresponse()\n",
" print(f\"Metering points API status: {res.status}\")\n",
"\n",
" if res.status != 200:\n",
" error_data = res.read().decode(\"utf-8\")\n",
" print(f\"Error: {error_data}\")\n",
" else:\n",
" data = res.read().decode(\"utf-8\")\n",
" print(f\"Response: {data[:50] + '...' if len(data) > 50 else data}\")"
]
},
{
"cell_type": "code",
"execution_count": 34,
"id": "686fe407",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Using metering point ID: 571313124600282119\n",
"Requesting consumption data from 2025-10-20 to 2025-10-27\n",
"Consumption data API status: 200\n",
"Success! Consumption data retrieved:\n",
"{\n",
" \"result\": [\n",
" {\n",
" \"MyEnergyData_MarketDocument\": {\n",
" \"mRID\": \"0HNGJHUMDF3GS:000001A6\",\n",
" \"createdDateTime\": \"2025-10-27T11:02:04Z\",\n",
" \"sender_MarketParticipant.name\": \"\",\n",
" \"sender_MarketParticipant.mRID\": {\n",
" \"codingScheme\": null,\n",
" \"name\": null\n",
" },\n",
" \"period.timeInterval\": {\n",
" \"start\": \"2025-10-19T22:00:00Z\",\n",
" \"end\": \"2025-10-26T23:00:00Z\"\n",
" },\n",
" \"TimeSeries\": [\n",
" {\n",
" \"mRID\": \"571313124600282119\",\n",
" \"businessType\": \"A04\",\n",
" \"curveType\": \"A01\",\n",
" \"measurement_Unit.name\": \"KWH\",\n",
" \"MarketEvaluationPoint\": {\n",
" \"mRID\": {\n",
" \"codingScheme\": \"A10\",\n",
" \"name\": \"571313124600282119\"\n",
" }\n",
" },\n",
" \"Period\": [\n",
" {\n",
" \"resolution\": \"PT1H\",\n",
" \"timeInterval\": {\n",
" \"start\": \"2025-10-19T22:00:00Z\",\n",
" \"end\": \"2025-10-2...\n",
"Consumption data API status: 200\n",
"Success! Consumption data retrieved:\n",
"{\n",
" \"result\": [\n",
" {\n",
" \"MyEnergyData_MarketDocument\": {\n",
" \"mRID\": \"0HNGJHUMDF3GS:000001A6\",\n",
" \"createdDateTime\": \"2025-10-27T11:02:04Z\",\n",
" \"sender_MarketParticipant.name\": \"\",\n",
" \"sender_MarketParticipant.mRID\": {\n",
" \"codingScheme\": null,\n",
" \"name\": null\n",
" },\n",
" \"period.timeInterval\": {\n",
" \"start\": \"2025-10-19T22:00:00Z\",\n",
" \"end\": \"2025-10-26T23:00:00Z\"\n",
" },\n",
" \"TimeSeries\": [\n",
" {\n",
" \"mRID\": \"571313124600282119\",\n",
" \"businessType\": \"A04\",\n",
" \"curveType\": \"A01\",\n",
" \"measurement_Unit.name\": \"KWH\",\n",
" \"MarketEvaluationPoint\": {\n",
" \"mRID\": {\n",
" \"codingScheme\": \"A10\",\n",
" \"name\": \"571313124600282119\"\n",
" }\n",
" },\n",
" \"Period\": [\n",
" {\n",
" \"resolution\": \"PT1H\",\n",
" \"timeInterval\": {\n",
" \"start\": \"2025-10-19T22:00:00Z\",\n",
" \"end\": \"2025-10-2...\n"
]
}
],
"source": [
"# Get consumption data for a specific metering point ID\n",
"\n",
"if not access_token:\n",
" print(\"No access token available\")\n",
"else:\n",
" # First, let's get the metering points data properly and extract the metering point ID\n",
" conn = http.client.HTTPSConnection(server_url.host)\n",
"\n",
" api_headers = {\n",
" \"Authorization\": f\"Bearer {access_token}\",\n",
" \"Content-Type\": \"application/json\",\n",
" \"api-version\": \"1.0\",\n",
" }\n",
"\n",
" # Get metering points\n",
" from urllib.parse import urlencode\n",
"\n",
" params = {\"includeAll\": True}\n",
"\n",
" conn.request(\n",
" \"GET\",\n",
" \"/customerapi/api/meteringpoints/meteringpoints?\" + urlencode(params),\n",
" headers=api_headers,\n",
" )\n",
" res = conn.getresponse()\n",
"\n",
" if res.status == 200:\n",
" metering_data = res.read().decode(\"utf-8\")\n",
" metering_points = json.loads(metering_data)\n",
"\n",
" if metering_points[\"result\"]:\n",
" # Get the first metering point ID\n",
" metering_point_id = metering_points[\"result\"][0][\"meteringPointId\"]\n",
" print(f\"Using metering point ID: {metering_point_id}\")\n",
"\n",
" # Now get consumption data for this metering point\n",
" # Let's get data from the last 7 days\n",
" from datetime import datetime, timedelta\n",
"\n",
" end_date = datetime.now()\n",
" start_date = end_date - timedelta(days=7)\n",
"\n",
" # Format dates as required by the API (YYYY-MM-DD)\n",
" date_from = start_date.strftime(\"%Y-%m-%d\")\n",
" date_to = end_date.strftime(\"%Y-%m-%d\")\n",
"\n",
" print(f\"Requesting consumption data from {date_from} to {date_to}\")\n",
"\n",
" # Prepare the request body for consumption data\n",
" consumption_request = {\n",
" \"meteringPoints\": {\"meteringPoint\": [metering_point_id]}\n",
" }\n",
"\n",
" # Create new connection for consumption data request\n",
" conn = http.client.HTTPSConnection(server_url.host)\n",
"\n",
" # Make the consumption data request\n",
" conn.request(\n",
" \"POST\",\n",
" f\"/customerapi/api/meterdata/gettimeseries/{date_from}/{date_to}/Hour\",\n",
" body=json.dumps(consumption_request),\n",
" headers=api_headers,\n",
" )\n",
"\n",
" res = conn.getresponse()\n",
" print(f\"Consumption data API status: {res.status}\")\n",
"\n",
" if res.status == 200:\n",
" consumption_data = res.read().decode(\"utf-8\")\n",
" consumption_json = json.loads(consumption_data)\n",
" print(\"Success! Consumption data retrieved:\")\n",
" print(\n",
" json.dumps(consumption_json, indent=2)[:1000] + \"...\"\n",
" if len(json.dumps(consumption_json, indent=2)) > 1000\n",
" else json.dumps(consumption_json, indent=2)\n",
" )\n",
" else:\n",
" error_data = res.read().decode(\"utf-8\")\n",
" print(f\"Error getting consumption data: {error_data}\")\n",
" else:\n",
" print(\"No metering points found\")\n",
" else:\n",
" error_data = res.read().decode(\"utf-8\")\n",
" print(f\"Error getting metering points: {error_data}\")"
]
},
{
"cell_type": "code",
"execution_count": 35,
"id": "97359b7b",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Metering Point ID: 571313124600282119\n",
"Unit: KWH\n",
"Data Period: 2025-10-19T22:00:00Z to 2025-10-26T23:00:00Z\n",
"\n",
"============================================================\n",
"CONSUMPTION DATA:\n",
"============================================================\n",
"\n",
"Period starting: 2025-10-19T22:00:00Z\n",
"Resolution: PT1H\n",
"Number of hourly readings: 24\n",
" Hour 1: 0.53 kWh (Quality: A04)\n",
" Hour 2: 0.55 kWh (Quality: A04)\n",
" Hour 3: 0.51 kWh (Quality: A04)\n",
" Hour 4: 0.5 kWh (Quality: A04)\n",
" Hour 5: 0.6 kWh (Quality: A04)\n",
" Hour 6: 0.52 kWh (Quality: A04)\n",
" Hour 7: 0.68 kWh (Quality: A04)\n",
" Hour 8: 0.69 kWh (Quality: A04)\n",
" Hour 9: 0.62 kWh (Quality: A04)\n",
" Hour 10: 0.62 kWh (Quality: A04)\n",
" ... and 14 more hourly readings\n",
"\n",
"Period starting: 2025-10-20T22:00:00Z\n",
"Resolution: PT1H\n",
"Number of hourly readings: 24\n",
" Hour 1: 0.58 kWh (Quality: A04)\n",
" Hour 2: 0.56 kWh (Quality: A04)\n",
" Hour 3: 0.57 kWh (Quality: A04)\n",
" Hour 4: 0.61 kWh (Quality: A04)\n",
" Hour 5: 0.64 kWh (Quality: A04)\n",
" Hour 6: 0.58 kWh (Quality: A04)\n",
" Hour 7: 0.62 kWh (Quality: A04)\n",
" Hour 8: 0.53 kWh (Quality: A04)\n",
" Hour 9: 0.58 kWh (Quality: A04)\n",
" Hour 10: 0.53 kWh (Quality: A04)\n",
" ... and 14 more hourly readings\n",
"\n",
"Period starting: 2025-10-21T22:00:00Z\n",
"Resolution: PT1H\n",
"Number of hourly readings: 24\n",
" Hour 1: 0.62 kWh (Quality: A04)\n",
" Hour 2: 0.62 kWh (Quality: A04)\n",
" Hour 3: 0.5 kWh (Quality: A04)\n",
" Hour 4: 0.59 kWh (Quality: A04)\n",
" Hour 5: 0.57 kWh (Quality: A04)\n",
" Hour 6: 0.58 kWh (Quality: A04)\n",
" Hour 7: 0.55 kWh (Quality: A04)\n",
" Hour 8: 0.69 kWh (Quality: A04)\n",
" Hour 9: 0.77 kWh (Quality: A04)\n",
" Hour 10: 1.0 kWh (Quality: A04)\n",
" ... and 14 more hourly readings\n",
"\n",
"Period starting: 2025-10-22T22:00:00Z\n",
"Resolution: PT1H\n",
"Number of hourly readings: 24\n",
" Hour 1: 0.59 kWh (Quality: A04)\n",
" Hour 2: 0.68 kWh (Quality: A04)\n",
" Hour 3: 0.6 kWh (Quality: A04)\n",
" Hour 4: 0.6 kWh (Quality: A04)\n",
" Hour 5: 0.57 kWh (Quality: A04)\n",
" Hour 6: 0.57 kWh (Quality: A04)\n",
" Hour 7: 0.57 kWh (Quality: A04)\n",
" Hour 8: 0.57 kWh (Quality: A04)\n",
" Hour 9: 0.59 kWh (Quality: A04)\n",
" Hour 10: 0.58 kWh (Quality: A04)\n",
" ... and 14 more hourly readings\n",
"\n",
"Period starting: 2025-10-23T22:00:00Z\n",
"Resolution: PT1H\n",
"Number of hourly readings: 24\n",
" Hour 1: 0.58 kWh (Quality: A04)\n",
" Hour 2: 0.72 kWh (Quality: A04)\n",
" Hour 3: 0.61 kWh (Quality: A04)\n",
" Hour 4: 0.6 kWh (Quality: A04)\n",
" Hour 5: 0.54 kWh (Quality: A04)\n",
" Hour 6: 0.53 kWh (Quality: A04)\n",
" Hour 7: 0.49 kWh (Quality: A04)\n",
" Hour 8: 0.8 kWh (Quality: A04)\n",
" Hour 9: 0.62 kWh (Quality: A04)\n",
" Hour 10: 1.21 kWh (Quality: A04)\n",
" ... and 14 more hourly readings\n",
"\n",
"Period starting: 2025-10-24T22:00:00Z\n",
"Resolution: PT1H\n",
"Number of hourly readings: 24\n",
" Hour 1: 0.89 kWh (Quality: A04)\n",
" Hour 2: 1.15 kWh (Quality: A04)\n",
" Hour 3: 0.58 kWh (Quality: A04)\n",
" Hour 4: 0.58 kWh (Quality: A04)\n",
" Hour 5: 0.62 kWh (Quality: A04)\n",
" Hour 6: 0.61 kWh (Quality: A04)\n",
" Hour 7: 0.58 kWh (Quality: A04)\n",
" Hour 8: 0.56 kWh (Quality: A04)\n",
" Hour 9: 0.53 kWh (Quality: A04)\n",
" Hour 10: 0.61 kWh (Quality: A04)\n",
" ... and 14 more hourly readings\n",
"\n",
"Period starting: 2025-10-25T22:00:00Z\n",
"Resolution: PT1H\n",
"Number of hourly readings: 25\n",
" Hour 1: 0.63 kWh (Quality: A04)\n",
" Hour 2: 0.51 kWh (Quality: A04)\n",
" Hour 3: 0.55 kWh (Quality: A04)\n",
" Hour 4: 0.56 kWh (Quality: A04)\n",
" Hour 5: 0.56 kWh (Quality: A04)\n",
" Hour 6: 0.5 kWh (Quality: A04)\n",
" Hour 7: 0.55 kWh (Quality: A04)\n",
" Hour 8: 0.56 kWh (Quality: A04)\n",
" Hour 9: 0.55 kWh (Quality: A04)\n",
" Hour 10: 0.54 kWh (Quality: A04)\n",
" ... and 15 more hourly readings\n",
"\n",
"============================================================\n",
"SUMMARY:\n",
"Total consumption over 169 hours: 138.31 kWh\n",
"Average hourly consumption: 0.818 kWh\n",
"============================================================\n"
]
}
],
"source": [
"# Parse and display the consumption data in a more readable format\n",
"\n",
"if \"consumption_json\" in locals() and consumption_json:\n",
" try:\n",
" # Extract the time series data\n",
" market_document = consumption_json[\"result\"][0][\"MyEnergyData_MarketDocument\"]\n",
" time_series = market_document[\"TimeSeries\"][0]\n",
" periods = time_series[\"Period\"]\n",
"\n",
" print(f\"Metering Point ID: {time_series['mRID']}\")\n",
" print(f\"Unit: {time_series['measurement_Unit.name']}\")\n",
" print(\n",
" f\"Data Period: {market_document['period.timeInterval']['start']} to {market_document['period.timeInterval']['end']}\"\n",
" )\n",
" print(\"\\n\" + \"=\" * 60)\n",
" print(\"CONSUMPTION DATA:\")\n",
" print(\"=\" * 60)\n",
"\n",
" total_consumption = 0\n",
" data_points = 0\n",
"\n",
" for period in periods:\n",
" period_start = period[\"timeInterval\"][\"start\"]\n",
" print(f\"\\nPeriod starting: {period_start}\")\n",
" print(f\"Resolution: {period['resolution']}\")\n",
"\n",
" if \"Point\" in period:\n",
" points = period[\"Point\"]\n",
" print(f\"Number of hourly readings: {len(points)}\")\n",
"\n",
" for point in points[:10]: # Show first 10 points\n",
" position = point[\"position\"]\n",
" quantity = float(point[\"out_Quantity.quantity\"])\n",
" quality = point[\"out_Quantity.quality\"]\n",
"\n",
" print(f\" Hour {position}: {quantity} kWh (Quality: {quality})\")\n",
" total_consumption += quantity\n",
" data_points += 1\n",
"\n",
" if len(points) > 10:\n",
" print(f\" ... and {len(points) - 10} more hourly readings\")\n",
" # Add remaining consumption to total\n",
" for point in points[10:]:\n",
" total_consumption += float(point[\"out_Quantity.quantity\"])\n",
" data_points += 1\n",
"\n",
" print(\"\\n\" + \"=\" * 60)\n",
" print(\"SUMMARY:\")\n",
" print(\n",
" f\"Total consumption over {data_points} hours: {total_consumption:.2f} kWh\"\n",
" )\n",
" print(f\"Average hourly consumption: {total_consumption / data_points:.3f} kWh\")\n",
" print(\"=\" * 60)\n",
"\n",
" except Exception as e:\n",
" print(f\"Error parsing consumption data: {e}\")\n",
" print(\"Raw data structure:\")\n",
" print(json.dumps(consumption_json, indent=2)[:500] + \"...\")\n",
"else:\n",
" print(\"No consumption data available to parse\")"
]
},
{
"cell_type": "markdown",
"id": "2bf461e8",
"metadata": {},
"source": [
"# PostgreSQL Database Design\n",
"\n",
"For storing energy consumption data efficiently, we'll design a normalized database schema with the following tables:\n",
"\n",
"## Tables Structure:\n",
"\n",
"### 1. `metering_points` - Store meter information\n",
"- `id` (Primary Key)\n",
"- `metering_point_id` (Unique identifier from API)\n",
"- `street_name`, `building_number`, `city` etc.\n",
"- `meter_number`\n",
"- `consumer_name`\n",
"- `created_at`, `updated_at`\n",
"\n",
"### 2. `consumption_readings` - Store hourly consumption data\n",
"- `id` (Primary Key)\n",
"- `metering_point_id` (Foreign Key)\n",
"- `timestamp` (UTC timestamp for the reading)\n",
"- `consumption_kwh` (Energy consumption in kWh)\n",
"- `quality` (Data quality indicator)\n",
"- `period_resolution` (e.g., 'PT1H' for hourly)\n",
"- `created_at`\n",
"\n",
"### Benefits:\n",
"- **Normalized**: Avoids data duplication\n",
"- **Indexed**: Fast queries on timestamp and metering_point_id\n",
"- **Scalable**: Can handle multiple meters and large time series data\n",
"- **Flexible**: Easy to add new meters or extend with additional fields"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "969c7953",
"metadata": {},
"outputs": [],
"source": [
"# SQL Schema for PostgreSQL Database\n",
"\n",
"create_tables_sql = \"\"\"\n",
"-- Create metering_points table\n",
"CREATE TABLE IF NOT EXISTS metering_points (\n",
" id SERIAL PRIMARY KEY,\n",
" metering_point_id VARCHAR(50) UNIQUE NOT NULL,\n",
" street_code VARCHAR(10),\n",
" street_name VARCHAR(255),\n",
" building_number VARCHAR(20),\n",
" floor_id VARCHAR(10),\n",
" room_id VARCHAR(10),\n",
" city_subdivision_name VARCHAR(255),\n",
" municipality_code VARCHAR(10),\n",
" location_description TEXT,\n",
" settlement_method VARCHAR(10),\n",
" meter_reading_occurrence VARCHAR(20),\n",
" first_consumer_party_name VARCHAR(255),\n",
" second_consumer_party_name VARCHAR(255),\n",
" meter_number VARCHAR(50),\n",
" consumer_start_date TIMESTAMP WITH TIME ZONE,\n",
" type_of_mp VARCHAR(10),\n",
" balance_supplier_name VARCHAR(255),\n",
" postcode VARCHAR(10),\n",
" created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,\n",
" updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP\n",
");\n",
"\n",
"-- Create consumption_readings table\n",
"CREATE TABLE IF NOT EXISTS consumption_readings (\n",
" id SERIAL PRIMARY KEY,\n",
" metering_point_id VARCHAR(50) REFERENCES metering_points(metering_point_id),\n",
" timestamp TIMESTAMP WITH TIME ZONE NOT NULL,\n",
" consumption_kwh DECIMAL(10, 3) NOT NULL,\n",
" quality VARCHAR(10),\n",
" period_resolution VARCHAR(10) DEFAULT 'PT1H',\n",
" created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,\n",
" UNIQUE(metering_point_id, timestamp)\n",
");\n",
"\n",
"-- Create indexes for performance\n",
"CREATE INDEX IF NOT EXISTS idx_consumption_metering_point ON consumption_readings(metering_point_id);\n",
"CREATE INDEX IF NOT EXISTS idx_consumption_timestamp ON consumption_readings(timestamp);\n",
"CREATE INDEX IF NOT EXISTS idx_consumption_metering_point_timestamp ON consumption_readings(metering_point_id, timestamp);\n",
"\n",
"-- Create a function to update the updated_at timestamp\n",
"CREATE OR REPLACE FUNCTION update_updated_at_column()\n",
"RETURNS TRIGGER AS $$\n",
"BEGIN\n",
" NEW.updated_at = CURRENT_TIMESTAMP;\n",
" RETURN NEW;\n",
"END;\n",
"$$ language 'plpgsql';\n",
"\n",
"-- Create trigger for metering_points\n",
"CREATE TRIGGER update_metering_points_updated_at \n",
" BEFORE UPDATE ON metering_points \n",
" FOR EACH ROW EXECUTE FUNCTION update_updated_at_column();\n",
"\"\"\"\n",
"\n",
"print(\"PostgreSQL Schema:\")\n",
"print(create_tables_sql)"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "b8159a41",
"metadata": {},
"outputs": [],
"source": [
"# Python code to store data in PostgreSQL using psycopg2\n",
"\n",
"import psycopg2\n",
"from datetime import datetime\n",
"\n",
"# Database connection configuration\n",
"DB_CONFIG = {\n",
" \"host\": \"localhost\",\n",
" \"database\": \"energy_consumption\",\n",
" \"user\": \"your_username\",\n",
" \"password\": \"your_password\",\n",
" \"port\": 5432,\n",
"}\n",
"\n",
"\n",
"class EnergyDataStorage:\n",
" def __init__(self, db_config):\n",
" self.db_config = db_config\n",
" self.connection = None\n",
"\n",
" def connect(self):\n",
" \"\"\"Establish database connection\"\"\"\n",
" try:\n",
" self.connection = psycopg2.connect(**self.db_config)\n",
" print(\"Database connection established\")\n",
" return True\n",
" except Exception as e:\n",
" print(f\"Error connecting to database: {e}\")\n",
" return False\n",
"\n",
" def create_tables(self):\n",
" \"\"\"Create database tables if they don't exist\"\"\"\n",
" if not self.connection:\n",
" print(\"No database connection\")\n",
" return False\n",
"\n",
" try:\n",
" with self.connection.cursor() as cursor:\n",
" cursor.execute(create_tables_sql)\n",
" self.connection.commit()\n",
" print(\"Tables created successfully\")\n",
" return True\n",
" except Exception as e:\n",
" print(f\"Error creating tables: {e}\")\n",
" self.connection.rollback()\n",
" return False\n",
"\n",
" def store_metering_point(self, metering_point_data):\n",
" \"\"\"Store or update metering point information\"\"\"\n",
" if not self.connection:\n",
" print(\"No database connection\")\n",
" return False\n",
"\n",
" try:\n",
" with self.connection.cursor() as cursor:\n",
" # Use UPSERT (INSERT ... ON CONFLICT)\n",
" upsert_sql = \"\"\"\n",
" INSERT INTO metering_points (\n",
" metering_point_id, street_code, street_name, building_number,\n",
" floor_id, room_id, city_subdivision_name, municipality_code,\n",
" location_description, settlement_method, meter_reading_occurrence,\n",
" first_consumer_party_name, second_consumer_party_name,\n",
" meter_number, consumer_start_date, type_of_mp,\n",
" balance_supplier_name, postcode\n",
" ) VALUES (\n",
" %(meteringPointId)s, %(streetCode)s, %(streetName)s, %(buildingNumber)s,\n",
" %(floorId)s, %(roomId)s, %(citySubDivisionName)s, %(municipalityCode)s,\n",
" %(locationDescription)s, %(settlementMethod)s, %(meterReadingOccurrence)s,\n",
" %(firstConsumerPartyName)s, %(secondConsumerPartyName)s,\n",
" %(meterNumber)s, %(consumerStartDate)s, %(typeOfMP)s,\n",
" %(balanceSupplierName)s, %(postcode)s\n",
" )\n",
" ON CONFLICT (metering_point_id) \n",
" DO UPDATE SET\n",
" street_code = EXCLUDED.street_code,\n",
" street_name = EXCLUDED.street_name,\n",
" building_number = EXCLUDED.building_number,\n",
" updated_at = CURRENT_TIMESTAMP;\n",
" \"\"\"\n",
"\n",
" cursor.execute(upsert_sql, metering_point_data)\n",
" self.connection.commit()\n",
" print(\n",
" f\"Metering point {metering_point_data['meteringPointId']} stored successfully\"\n",
" )\n",
" return True\n",
"\n",
" except Exception as e:\n",
" print(f\"Error storing metering point: {e}\")\n",
" self.connection.rollback()\n",
" return False\n",
"\n",
" def store_consumption_readings(self, metering_point_id, consumption_data):\n",
" \"\"\"Store consumption readings with conflict handling\"\"\"\n",
" if not self.connection:\n",
" print(\"No database connection\")\n",
" return False\n",
"\n",
" try:\n",
" with self.connection.cursor() as cursor:\n",
" # Prepare batch insert with conflict handling\n",
" insert_sql = \"\"\"\n",
" INSERT INTO consumption_readings (\n",
" metering_point_id, timestamp, consumption_kwh, quality, period_resolution\n",
" ) VALUES (\n",
" %s, %s, %s, %s, %s\n",
" )\n",
" ON CONFLICT (metering_point_id, timestamp) \n",
" DO UPDATE SET\n",
" consumption_kwh = EXCLUDED.consumption_kwh,\n",
" quality = EXCLUDED.quality;\n",
" \"\"\"\n",
"\n",
" readings_data = []\n",
" for reading in consumption_data:\n",
" readings_data.append(\n",
" (\n",
" metering_point_id,\n",
" reading[\"timestamp\"],\n",
" reading[\"consumption_kwh\"],\n",
" reading[\"quality\"],\n",
" reading.get(\"period_resolution\", \"PT1H\"),\n",
" )\n",
" )\n",
"\n",
" cursor.executemany(insert_sql, readings_data)\n",
" self.connection.commit()\n",
" print(f\"Stored {len(readings_data)} consumption readings\")\n",
" return True\n",
"\n",
" except Exception as e:\n",
" print(f\"Error storing consumption readings: {e}\")\n",
" self.connection.rollback()\n",
" return False\n",
"\n",
" def close(self):\n",
" \"\"\"Close database connection\"\"\"\n",
" if self.connection:\n",
" self.connection.close()\n",
" print(\"Database connection closed\")\n",
"\n",
"\n",
"# Example usage:\n",
"print(\"EnergyDataStorage class defined. Use it like this:\")\n",
"print(\"\"\"\n",
"# Initialize storage\n",
"storage = EnergyDataStorage(DB_CONFIG)\n",
"storage.connect()\n",
"storage.create_tables()\n",
"\n",
"# Store metering point data\n",
"storage.store_metering_point(metering_point_data)\n",
"\n",
"# Store consumption readings\n",
"storage.store_consumption_readings(metering_point_id, readings_list)\n",
"\n",
"storage.close()\n",
"\"\"\")"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "d00bdc8f",
"metadata": {},
"outputs": [],
"source": [
"# Function to transform API data for database storage\n",
"\n",
"\n",
"def transform_consumption_data_for_db(consumption_json):\n",
" \"\"\"\n",
" Transform the Eloverblik API response into database-ready format\n",
" \"\"\"\n",
" if not consumption_json or \"result\" not in consumption_json:\n",
" return None, []\n",
"\n",
" try:\n",
" # Extract metering point data\n",
" market_document = consumption_json[\"result\"][0][\"MyEnergyData_MarketDocument\"]\n",
" time_series = market_document[\"TimeSeries\"][0]\n",
" metering_point_id = time_series[\"mRID\"]\n",
"\n",
" # For demonstration, we'll create a basic metering point record\n",
" # In practice, you'd get this from the metering points API call\n",
" metering_point_data = {\n",
" \"meteringPointId\": metering_point_id,\n",
" \"streetCode\": None,\n",
" \"streetName\": None,\n",
" \"buildingNumber\": None,\n",
" \"floorId\": None,\n",
" \"roomId\": None,\n",
" \"citySubDivisionName\": None,\n",
" \"municipalityCode\": None,\n",
" \"locationDescription\": None,\n",
" \"settlementMethod\": None,\n",
" \"meterReadingOccurrence\": None,\n",
" \"firstConsumerPartyName\": None,\n",
" \"secondConsumerPartyName\": None,\n",
" \"meterNumber\": None,\n",
" \"consumerStartDate\": None,\n",
" \"typeOfMP\": None,\n",
" \"balanceSupplierName\": None,\n",
" \"postcode\": None,\n",
" }\n",
"\n",
" # Extract consumption readings\n",
" consumption_readings = []\n",
" periods = time_series[\"Period\"]\n",
"\n",
" for period in periods:\n",
" period_start = datetime.fromisoformat(\n",
" period[\"timeInterval\"][\"start\"].replace(\"Z\", \"+00:00\")\n",
" )\n",
" resolution = period[\"resolution\"]\n",
"\n",
" if \"Point\" in period:\n",
" for point in period[\"Point\"]:\n",
" # Calculate the actual timestamp for this point\n",
" position = int(point[\"position\"])\n",
" # Position is 1-based, so subtract 1 to get hours offset\n",
" hours_offset = position - 1\n",
"\n",
" point_timestamp = period_start + timedelta(hours=hours_offset)\n",
"\n",
" reading = {\n",
" \"timestamp\": point_timestamp,\n",
" \"consumption_kwh\": float(point[\"out_Quantity.quantity\"]),\n",
" \"quality\": point[\"out_Quantity.quality\"],\n",
" \"period_resolution\": resolution,\n",
" }\n",
" consumption_readings.append(reading)\n",
"\n",
" return metering_point_data, consumption_readings\n",
"\n",
" except Exception as e:\n",
" print(f\"Error transforming data: {e}\")\n",
" return None, []\n",
"\n",
"\n",
"# Test the transformation with our existing data\n",
"if \"consumption_json\" in locals() and consumption_json:\n",
" metering_point, readings = transform_consumption_data_for_db(consumption_json)\n",
"\n",
" if metering_point and readings:\n",
" print(\n",
" f\"Transformed data for metering point: {metering_point['meteringPointId']}\"\n",
" )\n",
" print(f\"Number of readings: {len(readings)}\")\n",
" print(\"\\nSample readings:\")\n",
" for i, reading in enumerate(readings[:3]):\n",
" print(\n",
" f\" {i + 1}. {reading['timestamp']}: {reading['consumption_kwh']} kWh (Quality: {reading['quality']})\"\n",
" )\n",
" print(\" ...\")\n",
" else:\n",
" print(\"Failed to transform data\")\n",
"else:\n",
" print(\"No consumption data available to transform\")"
]
}
],
"metadata": {
"kernelspec": {
"display_name": "energy-consumption-ingester",
"language": "python",
"name": "python3"
},
"language_info": {
"codemirror_mode": {
"name": "ipython",
"version": 3
},
"file_extension": ".py",
"mimetype": "text/x-python",
"name": "python",
"nbconvert_exporter": "python",
"pygments_lexer": "ipython3",
"version": "3.12.1"
}
},
"nbformat": 4,
"nbformat_minor": 5
}