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

39 KiB

Setup

In [ ]:
import os
from dotenv import load_dotenv
import http.client
import json
from pydantic import AnyUrl

load_dotenv()
ELOVERBLIK_API_TOKEN = os.getenv("ELOVERBLIK_API_TOKEN")


server_url = AnyUrl("https://api.eloverblik.dk")

# Prepare JWT authentication

conn = http.client.HTTPSConnection(server_url.host)

headers = {
    "Authorization": f"Bearer {ELOVERBLIK_API_TOKEN}",
    "Content-Type": "application/json",
    "api-version": "1.0",
}

conn.request("GET", "/customerapi/api/token", headers=headers)

res = conn.getresponse()
print(f"Token request status: {res.status}")

if res.status == 200:
    token_response = res.read().decode("utf-8")
    print(f"Token response: {token_response[:50]}...")

    # Parse the response to get the access token
    try:
        token_data = json.loads(token_response)
        access_token = token_data["result"]
        print(f"Successfully obtained access token: {access_token[:50]}...")
    except Exception as e:
        print(f"Error parsing token response: {e}")
        access_token = None
else:
    error_data = res.read().decode("utf-8")
    print(f"Error getting token: {error_data}")
    access_token = None
Token request status: 200
Token response: {"result":"eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJ0b2tlblR5cGUiOiJDdXN0b21lckFQSV9EYXRhQWNjZXNzIiwidG9rZW5pZCI6ImE2NzY1ODRmLWY1MzktNDI0Yy04MmZhLWM5MjNhMjNhOTVlYSIsIndlYkFwcCI6IkN1c3RvbWVyQXBwIiwidmVyc2lvbiI6IjIiLCJpZGVudGl0eVRva2VuIjoiWUpoaXlqMXRjMUZiTnBtWmNaM0tjdGUrbzNzQ1FxR3B1dWtEcVh4QzVveEhKMisvcTAycTdpWGlZZURTRm9jUjVmYVBZMHVCeGpiWnYrOE52SHUycnpFU2pFRDNUajBQbmM5QzBZY1FXb0hzTDRvVHliOWRTSysrVFR5NTdDWTR4bit1dTdYN2lYdTNyV1pLUTUzTkh5T2NRSjB3Y3ZtWURHUjVDT3lITW14WDZhMjFEUHdBOWFsRk53RlFXbGkxT0NuTnp3QUdQaXJZWnZLcTYxNW52aFBjczZaN0FWUlh5MnRqMTFPbWpuR1NQVFJYU2pvdjdpRkhZcmhTdTRiTFpDQXhDd0FHWHl4a0tPWENRRDZFUVdTUU92d0ZpMHdnSUZwVHBIQjJKaHlEVFRvM2FYU1RaSnh2MEpFZktvalNLWXUxSE54enpWcTE3NmVWei9UVW5JZTdUeTc0LzZNekFqZmplcWpFUTk3aXcvbnRwTFNzdGkwakZqdjB6V1hKRXRKMVNQYURmbG5ZdG9DcEtnZmdCVDc4RTNKcXVIZSszQjJLanUwL0xlUU4rMXpHU2cxeEhuUFZOdlcxanRabjhWbjhIaVVmNUhUK0tPZ3ZDNitqaWJQRWRwNXJvaWtSNkFjRzZHcS8xOThkL2d6bEo4MzJMVWVnR1BiTlBHVVBwZ2VQTmFtUlRSWXpnSk9GdmwzV0NGb1NXMUZhYU5UalBJL1F4Wjh4L05VVzlJZXFXbG1rb1VrMVhIS3lWSzBjdW1KWXlaelJBNHJKdGJ0QXVaaGdMUCtpUmtSUzlFMkcwOWhDQnhzOHJKWjRycWVvZmVITC9iQkQwYXIzd25tREFpUks3L1ovNmkvbXFJMEZHa3NUcXdYZzdBTTY0SldGQkNUcm1ZMk9hSlIwRHpSZmd3ekN5Y0NzRHRRUW1RNTZwNWVPRzJyRnlUUHh5YjEwS3p4a3FTbTZVNDExcEhlVllFbnR4cVloYUg4and2RWJxVFJBNEZ5VUZtQVVxQmhIVDFWNWpsRkFZL0U3bjUwejAxSE5WL21KNzVBcGNqektsaUx4YVBuQXdYUG43NmVVbHc4NUt4anJueEh5d3V6a00wVFhDSVRIZ0NrWmdYbkdxVEYwZE1TL3VhRE9Tbjl5NTY5SE9vb1FxM3VNZ0prV2RFaUdSMkUvcVErbFVNS05kVlNEL0F6Y0tLZkdhdkZUTko2WVhZOVozKzZKN2lpMWFDbGxoWFAwQ1VRRzhaOVpaMEMrU3JlN3FNSzVvaEpwY2FDOUdoZ1hVNklqZ0xZd2JTM0RGOUN3NGt3cmxHRVJSNURVZVNXREg2aFp5SVdEWWc2d1k4bHMxNGVlVjQ4bW5aTVVCSVE4cDZvNzhGSzFLUWc2V3BQZ3NqMGVOcXdBVW1HemNQcUs5UHJsdDc5VmZvK2pQSDlkUW1IaUJIbjFDcjROWWJpUEI4bllNN09JdFJpb3lHVmpyTi9lVjR1SGloc1c1ZzUxMG9pMjVxNzU0Y2VPbW10ck0ra0VzRGtoVEd1UzNxZml5ZUhqRDV0dk1leXU3eGs0c1gyeDdhVkdDeGV0eGlDalNmWkZ0R2JiWmY3VWhKWTAzNzBZbDYvZFRNeDJ4VFNlUzhqOUgwZUZJT1RBL3N3ZVNXUjdDMlc2MitWOC9SNERZT3lOdjJzWmRib2RIN3hKdk0yUnNxcHFRMUtqS1VsclpEUmRTMFk4ZDZTOGJuTTBGUC9tRXlDZkF1MmthRnE5d0hQcGZUdzJ2QVFIeTdMNnR5ZjNoSUNtOHQ1eGh4NEJwV1ZUUkVtQUVXZk5KQ3c2QXZlUFExVHY5TjVmQ1N6NXo5RlU2VmhxcUlQL1hWM2l3dUFQc3dVc2hpaG90ZTBPR2YrZlpyMEZHY04vSTJXYzZneDdIT3RYVitUWEdieGdvY2lianFWRTFyQTRweTJsYnpSQlo3dXVRc29OcW1YQlVkckF6U0V3MkFEZGg3cE1KWGpBRXF0cHZKOXluWUlqWWRpZXdocllEazRRNCtEekl5V01NRk5YOGVyakFER21PYnJvazB6aUtGd1Fwb3hCL2I2SU9GYzVwWjFvUm0zY0dNUjU2M3BPdmw0aTNIalhMbmxvbWkxQ29Hb1JPL203MGp2RXRMdXBvOGRRSGpFZXlEUzBmeGl1b0hHRWV1Uk12SGpjcVUranBadjloUjEwMXh1QVVGRjNaUFJ6NnhQb2JsZVR0N0xnSjZWMm9KZmJ4QmpPWEdLTEthdGNiUEkzMS84MWtSWlIyN0w5eVZ4NkRLays1ZjV4Zk5JdUxWd0pzeDNBRVBob2o4SERPVFRKSERMUlZpa1hOdGE4OTNZOWVaZ2xpRzdyMXZCaXRTZ2dwQTRlbDhjMHlIdGJDaG9Vb3N6SmsyMHZJWGphem1zRG5kTU5DVXluRGlWbEN4VDBROWRPMlc4Rlp5NERxRmxNN1BQVVZWZzRhMkRIUzhianZQR3dWZ0R6WDkrVzA3VEMrVVlpZ3k2Sk93YnkwMThqd2loQjlLb3d6M0Vqc1djTjF1RGRsYXZPbHNWbHFBU05tWUF0MmpnbVF1YnQzTUlrdFhObGNOMWl0MnA5NFJ3UE9qNmhjYm1OM3ZRaXNvVExyd1llK1RlciIsImh0dHA6Ly9zY2hlbWFzLnhtbHNvYXAub3JnL3dzLzIwMDUvMDUvaWRlbnRpdHkvY2xhaW1zL25hbWVpZGVudGlmaWVyIjoiTUlEOmUyZjE2NDE4LTlmZDItNGFkZi1hNDA2LWEwMDQ5ZTZjYzU4YSIsImh0dHA6Ly9zY2hlbWFzLnhtbHNvYXAub3JnL3dzLzIwMDUvMDUvaWRlbnRpdHkvY2xhaW1zL2dpdmVubmFtZSI6IkJyaWFuIEJqYXJrZSBKZW5zZW4iLCJsb2dpblR5cGUiOiJLZXlDYXJkIiwiYjNmIjoiRjRSQ1VHSjdoTUUyMlpUVzRnTGZha2crL0VXL2Y4MmZVVTFoNVFjMk1DWT0iLCJwaWQiOiJQSUQ6OTIwOC0yMDAyLTItNjAwNTA4Mzg2NzY4IiwidXNlcklkIjoiMTA0MzQ1OSIsImV4cCI6MTc2MTY0NzgxNywiaXNzIjoiRW5lcmdpbmV0IiwianRpIjoiYTY3NjU4NGYtZjUzOS00MjRjLTgyZmEtYzkyM2EyM2E5NWVhIiwidG9rZW5OYW1lIjoiZ2l0ZWEtYWNjZXNzLXRva2VuIiwiYXVkIjoiRW5lcmdpbmV0In0.zmb-Oy1Jn6u16dAodmdd3OApFbUVcTjetM5Z8x_evxQ"}
Successfully obtained access token: eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJ0b2tlblR5c...
In [33]:
if not access_token:
    print("No access token available to test with")
else:
    print("Testing access token with metering points API...")

    # Create new connection for API calls
    conn = http.client.HTTPSConnection(server_url.host)

    # Use the access token for API calls
    api_headers = {
        "Authorization": f"Bearer {access_token}",
        "Content-Type": "application/json",
        "api-version": "1.0",
    }

    # Get metering points
    from urllib.parse import urlencode

    params = {
        "includeAll": True,
    }

    conn.request(
        "GET",
        "/customerapi/api/meteringpoints/meteringpoints?" + urlencode(params),
        headers=api_headers,
    )

    res = conn.getresponse()
    print(f"Metering points API status: {res.status}")

    if res.status != 200:
        error_data = res.read().decode("utf-8")
        print(f"Error: {error_data}")
    else:
        data = res.read().decode("utf-8")
        print(f"Response: {data[:50] + '...' if len(data) > 50 else data}")
Testing access token with metering points API...
Metering points API status: 200
Response: {"result":[{"streetCode":"2724","streetName":"Vejl...
Metering points API status: 200
Response: {"result":[{"streetCode":"2724","streetName":"Vejl...
In [34]:
# Get consumption data for a specific metering point ID

if not access_token:
    print("No access token available")
else:
    # First, let's get the metering points data properly and extract the metering point ID
    conn = http.client.HTTPSConnection(server_url.host)

    api_headers = {
        "Authorization": f"Bearer {access_token}",
        "Content-Type": "application/json",
        "api-version": "1.0",
    }

    # Get metering points
    from urllib.parse import urlencode

    params = {"includeAll": True}

    conn.request(
        "GET",
        "/customerapi/api/meteringpoints/meteringpoints?" + urlencode(params),
        headers=api_headers,
    )
    res = conn.getresponse()

    if res.status == 200:
        metering_data = res.read().decode("utf-8")
        metering_points = json.loads(metering_data)

        if metering_points["result"]:
            # Get the first metering point ID
            metering_point_id = metering_points["result"][0]["meteringPointId"]
            print(f"Using metering point ID: {metering_point_id}")

            # Now get consumption data for this metering point
            # Let's get data from the last 7 days
            from datetime import datetime, timedelta

            end_date = datetime.now()
            start_date = end_date - timedelta(days=7)

            # Format dates as required by the API (YYYY-MM-DD)
            date_from = start_date.strftime("%Y-%m-%d")
            date_to = end_date.strftime("%Y-%m-%d")

            print(f"Requesting consumption data from {date_from} to {date_to}")

            # Prepare the request body for consumption data
            consumption_request = {
                "meteringPoints": {"meteringPoint": [metering_point_id]}
            }

            # Create new connection for consumption data request
            conn = http.client.HTTPSConnection(server_url.host)

            # Make the consumption data request
            conn.request(
                "POST",
                f"/customerapi/api/meterdata/gettimeseries/{date_from}/{date_to}/Hour",
                body=json.dumps(consumption_request),
                headers=api_headers,
            )

            res = conn.getresponse()
            print(f"Consumption data API status: {res.status}")

            if res.status == 200:
                consumption_data = res.read().decode("utf-8")
                consumption_json = json.loads(consumption_data)
                print("Success! Consumption data retrieved:")
                print(
                    json.dumps(consumption_json, indent=2)[:1000] + "..."
                    if len(json.dumps(consumption_json, indent=2)) > 1000
                    else json.dumps(consumption_json, indent=2)
                )
            else:
                error_data = res.read().decode("utf-8")
                print(f"Error getting consumption data: {error_data}")
        else:
            print("No metering points found")
    else:
        error_data = res.read().decode("utf-8")
        print(f"Error getting metering points: {error_data}")
Using metering point ID: 571313124600282119
Requesting consumption data from 2025-10-20 to 2025-10-27
Consumption data API status: 200
Success! Consumption data retrieved:
{
  "result": [
    {
      "MyEnergyData_MarketDocument": {
        "mRID": "0HNGJHUMDF3GS:000001A6",
        "createdDateTime": "2025-10-27T11:02:04Z",
        "sender_MarketParticipant.name": "",
        "sender_MarketParticipant.mRID": {
          "codingScheme": null,
          "name": null
        },
        "period.timeInterval": {
          "start": "2025-10-19T22:00:00Z",
          "end": "2025-10-26T23:00:00Z"
        },
        "TimeSeries": [
          {
            "mRID": "571313124600282119",
            "businessType": "A04",
            "curveType": "A01",
            "measurement_Unit.name": "KWH",
            "MarketEvaluationPoint": {
              "mRID": {
                "codingScheme": "A10",
                "name": "571313124600282119"
              }
            },
            "Period": [
              {
                "resolution": "PT1H",
                "timeInterval": {
                  "start": "2025-10-19T22:00:00Z",
                  "end": "2025-10-2...
Consumption data API status: 200
Success! Consumption data retrieved:
{
  "result": [
    {
      "MyEnergyData_MarketDocument": {
        "mRID": "0HNGJHUMDF3GS:000001A6",
        "createdDateTime": "2025-10-27T11:02:04Z",
        "sender_MarketParticipant.name": "",
        "sender_MarketParticipant.mRID": {
          "codingScheme": null,
          "name": null
        },
        "period.timeInterval": {
          "start": "2025-10-19T22:00:00Z",
          "end": "2025-10-26T23:00:00Z"
        },
        "TimeSeries": [
          {
            "mRID": "571313124600282119",
            "businessType": "A04",
            "curveType": "A01",
            "measurement_Unit.name": "KWH",
            "MarketEvaluationPoint": {
              "mRID": {
                "codingScheme": "A10",
                "name": "571313124600282119"
              }
            },
            "Period": [
              {
                "resolution": "PT1H",
                "timeInterval": {
                  "start": "2025-10-19T22:00:00Z",
                  "end": "2025-10-2...
In [35]:
# Parse and display the consumption data in a more readable format

if "consumption_json" in locals() and consumption_json:
    try:
        # Extract the time series data
        market_document = consumption_json["result"][0]["MyEnergyData_MarketDocument"]
        time_series = market_document["TimeSeries"][0]
        periods = time_series["Period"]

        print(f"Metering Point ID: {time_series['mRID']}")
        print(f"Unit: {time_series['measurement_Unit.name']}")
        print(
            f"Data Period: {market_document['period.timeInterval']['start']} to {market_document['period.timeInterval']['end']}"
        )
        print("\n" + "=" * 60)
        print("CONSUMPTION DATA:")
        print("=" * 60)

        total_consumption = 0
        data_points = 0

        for period in periods:
            period_start = period["timeInterval"]["start"]
            print(f"\nPeriod starting: {period_start}")
            print(f"Resolution: {period['resolution']}")

            if "Point" in period:
                points = period["Point"]
                print(f"Number of hourly readings: {len(points)}")

                for point in points[:10]:  # Show first 10 points
                    position = point["position"]
                    quantity = float(point["out_Quantity.quantity"])
                    quality = point["out_Quantity.quality"]

                    print(f"  Hour {position}: {quantity} kWh (Quality: {quality})")
                    total_consumption += quantity
                    data_points += 1

                if len(points) > 10:
                    print(f"  ... and {len(points) - 10} more hourly readings")
                    # Add remaining consumption to total
                    for point in points[10:]:
                        total_consumption += float(point["out_Quantity.quantity"])
                        data_points += 1

        print("\n" + "=" * 60)
        print("SUMMARY:")
        print(
            f"Total consumption over {data_points} hours: {total_consumption:.2f} kWh"
        )
        print(f"Average hourly consumption: {total_consumption / data_points:.3f} kWh")
        print("=" * 60)

    except Exception as e:
        print(f"Error parsing consumption data: {e}")
        print("Raw data structure:")
        print(json.dumps(consumption_json, indent=2)[:500] + "...")
else:
    print("No consumption data available to parse")
Metering Point ID: 571313124600282119
Unit: KWH
Data Period: 2025-10-19T22:00:00Z to 2025-10-26T23:00:00Z

============================================================
CONSUMPTION DATA:
============================================================

Period starting: 2025-10-19T22:00:00Z
Resolution: PT1H
Number of hourly readings: 24
  Hour 1: 0.53 kWh (Quality: A04)
  Hour 2: 0.55 kWh (Quality: A04)
  Hour 3: 0.51 kWh (Quality: A04)
  Hour 4: 0.5 kWh (Quality: A04)
  Hour 5: 0.6 kWh (Quality: A04)
  Hour 6: 0.52 kWh (Quality: A04)
  Hour 7: 0.68 kWh (Quality: A04)
  Hour 8: 0.69 kWh (Quality: A04)
  Hour 9: 0.62 kWh (Quality: A04)
  Hour 10: 0.62 kWh (Quality: A04)
  ... and 14 more hourly readings

Period starting: 2025-10-20T22:00:00Z
Resolution: PT1H
Number of hourly readings: 24
  Hour 1: 0.58 kWh (Quality: A04)
  Hour 2: 0.56 kWh (Quality: A04)
  Hour 3: 0.57 kWh (Quality: A04)
  Hour 4: 0.61 kWh (Quality: A04)
  Hour 5: 0.64 kWh (Quality: A04)
  Hour 6: 0.58 kWh (Quality: A04)
  Hour 7: 0.62 kWh (Quality: A04)
  Hour 8: 0.53 kWh (Quality: A04)
  Hour 9: 0.58 kWh (Quality: A04)
  Hour 10: 0.53 kWh (Quality: A04)
  ... and 14 more hourly readings

Period starting: 2025-10-21T22:00:00Z
Resolution: PT1H
Number of hourly readings: 24
  Hour 1: 0.62 kWh (Quality: A04)
  Hour 2: 0.62 kWh (Quality: A04)
  Hour 3: 0.5 kWh (Quality: A04)
  Hour 4: 0.59 kWh (Quality: A04)
  Hour 5: 0.57 kWh (Quality: A04)
  Hour 6: 0.58 kWh (Quality: A04)
  Hour 7: 0.55 kWh (Quality: A04)
  Hour 8: 0.69 kWh (Quality: A04)
  Hour 9: 0.77 kWh (Quality: A04)
  Hour 10: 1.0 kWh (Quality: A04)
  ... and 14 more hourly readings

Period starting: 2025-10-22T22:00:00Z
Resolution: PT1H
Number of hourly readings: 24
  Hour 1: 0.59 kWh (Quality: A04)
  Hour 2: 0.68 kWh (Quality: A04)
  Hour 3: 0.6 kWh (Quality: A04)
  Hour 4: 0.6 kWh (Quality: A04)
  Hour 5: 0.57 kWh (Quality: A04)
  Hour 6: 0.57 kWh (Quality: A04)
  Hour 7: 0.57 kWh (Quality: A04)
  Hour 8: 0.57 kWh (Quality: A04)
  Hour 9: 0.59 kWh (Quality: A04)
  Hour 10: 0.58 kWh (Quality: A04)
  ... and 14 more hourly readings

Period starting: 2025-10-23T22:00:00Z
Resolution: PT1H
Number of hourly readings: 24
  Hour 1: 0.58 kWh (Quality: A04)
  Hour 2: 0.72 kWh (Quality: A04)
  Hour 3: 0.61 kWh (Quality: A04)
  Hour 4: 0.6 kWh (Quality: A04)
  Hour 5: 0.54 kWh (Quality: A04)
  Hour 6: 0.53 kWh (Quality: A04)
  Hour 7: 0.49 kWh (Quality: A04)
  Hour 8: 0.8 kWh (Quality: A04)
  Hour 9: 0.62 kWh (Quality: A04)
  Hour 10: 1.21 kWh (Quality: A04)
  ... and 14 more hourly readings

Period starting: 2025-10-24T22:00:00Z
Resolution: PT1H
Number of hourly readings: 24
  Hour 1: 0.89 kWh (Quality: A04)
  Hour 2: 1.15 kWh (Quality: A04)
  Hour 3: 0.58 kWh (Quality: A04)
  Hour 4: 0.58 kWh (Quality: A04)
  Hour 5: 0.62 kWh (Quality: A04)
  Hour 6: 0.61 kWh (Quality: A04)
  Hour 7: 0.58 kWh (Quality: A04)
  Hour 8: 0.56 kWh (Quality: A04)
  Hour 9: 0.53 kWh (Quality: A04)
  Hour 10: 0.61 kWh (Quality: A04)
  ... and 14 more hourly readings

Period starting: 2025-10-25T22:00:00Z
Resolution: PT1H
Number of hourly readings: 25
  Hour 1: 0.63 kWh (Quality: A04)
  Hour 2: 0.51 kWh (Quality: A04)
  Hour 3: 0.55 kWh (Quality: A04)
  Hour 4: 0.56 kWh (Quality: A04)
  Hour 5: 0.56 kWh (Quality: A04)
  Hour 6: 0.5 kWh (Quality: A04)
  Hour 7: 0.55 kWh (Quality: A04)
  Hour 8: 0.56 kWh (Quality: A04)
  Hour 9: 0.55 kWh (Quality: A04)
  Hour 10: 0.54 kWh (Quality: A04)
  ... and 15 more hourly readings

============================================================
SUMMARY:
Total consumption over 169 hours: 138.31 kWh
Average hourly consumption: 0.818 kWh
============================================================

PostgreSQL Database Design

For storing energy consumption data efficiently, we'll design a normalized database schema with the following tables:

Tables Structure:

1. metering_points - Store meter information

  • id (Primary Key)
  • metering_point_id (Unique identifier from API)
  • street_name, building_number, city etc.
  • meter_number
  • consumer_name
  • created_at, updated_at

2. consumption_readings - Store hourly consumption data

  • id (Primary Key)
  • metering_point_id (Foreign Key)
  • timestamp (UTC timestamp for the reading)
  • consumption_kwh (Energy consumption in kWh)
  • quality (Data quality indicator)
  • period_resolution (e.g., 'PT1H' for hourly)
  • created_at

Benefits:

  • Normalized: Avoids data duplication
  • Indexed: Fast queries on timestamp and metering_point_id
  • Scalable: Can handle multiple meters and large time series data
  • Flexible: Easy to add new meters or extend with additional fields
In [ ]:
# SQL Schema for PostgreSQL Database

create_tables_sql = """
-- Create metering_points table
CREATE TABLE IF NOT EXISTS metering_points (
    id SERIAL PRIMARY KEY,
    metering_point_id VARCHAR(50) UNIQUE NOT NULL,
    street_code VARCHAR(10),
    street_name VARCHAR(255),
    building_number VARCHAR(20),
    floor_id VARCHAR(10),
    room_id VARCHAR(10),
    city_subdivision_name VARCHAR(255),
    municipality_code VARCHAR(10),
    location_description TEXT,
    settlement_method VARCHAR(10),
    meter_reading_occurrence VARCHAR(20),
    first_consumer_party_name VARCHAR(255),
    second_consumer_party_name VARCHAR(255),
    meter_number VARCHAR(50),
    consumer_start_date TIMESTAMP WITH TIME ZONE,
    type_of_mp VARCHAR(10),
    balance_supplier_name VARCHAR(255),
    postcode VARCHAR(10),
    created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
    updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

-- Create consumption_readings table
CREATE TABLE IF NOT EXISTS consumption_readings (
    id SERIAL PRIMARY KEY,
    metering_point_id VARCHAR(50) REFERENCES metering_points(metering_point_id),
    timestamp TIMESTAMP WITH TIME ZONE NOT NULL,
    consumption_kwh DECIMAL(10, 3) NOT NULL,
    quality VARCHAR(10),
    period_resolution VARCHAR(10) DEFAULT 'PT1H',
    created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
    UNIQUE(metering_point_id, timestamp)
);

-- Create indexes for performance
CREATE INDEX IF NOT EXISTS idx_consumption_metering_point ON consumption_readings(metering_point_id);
CREATE INDEX IF NOT EXISTS idx_consumption_timestamp ON consumption_readings(timestamp);
CREATE INDEX IF NOT EXISTS idx_consumption_metering_point_timestamp ON consumption_readings(metering_point_id, timestamp);

-- Create a function to update the updated_at timestamp
CREATE OR REPLACE FUNCTION update_updated_at_column()
RETURNS TRIGGER AS $$
BEGIN
    NEW.updated_at = CURRENT_TIMESTAMP;
    RETURN NEW;
END;
$$ language 'plpgsql';

-- Create trigger for metering_points
CREATE TRIGGER update_metering_points_updated_at 
    BEFORE UPDATE ON metering_points 
    FOR EACH ROW EXECUTE FUNCTION update_updated_at_column();
"""

print("PostgreSQL Schema:")
print(create_tables_sql)
In [ ]:
# Python code to store data in PostgreSQL using psycopg2

import psycopg2
from datetime import datetime

# Database connection configuration
DB_CONFIG = {
    "host": "localhost",
    "database": "energy_consumption",
    "user": "your_username",
    "password": "your_password",
    "port": 5432,
}


class EnergyDataStorage:
    def __init__(self, db_config):
        self.db_config = db_config
        self.connection = None

    def connect(self):
        """Establish database connection"""
        try:
            self.connection = psycopg2.connect(**self.db_config)
            print("Database connection established")
            return True
        except Exception as e:
            print(f"Error connecting to database: {e}")
            return False

    def create_tables(self):
        """Create database tables if they don't exist"""
        if not self.connection:
            print("No database connection")
            return False

        try:
            with self.connection.cursor() as cursor:
                cursor.execute(create_tables_sql)
                self.connection.commit()
                print("Tables created successfully")
                return True
        except Exception as e:
            print(f"Error creating tables: {e}")
            self.connection.rollback()
            return False

    def store_metering_point(self, metering_point_data):
        """Store or update metering point information"""
        if not self.connection:
            print("No database connection")
            return False

        try:
            with self.connection.cursor() as cursor:
                # Use UPSERT (INSERT ... ON CONFLICT)
                upsert_sql = """
                INSERT INTO metering_points (
                    metering_point_id, street_code, street_name, building_number,
                    floor_id, room_id, city_subdivision_name, municipality_code,
                    location_description, settlement_method, meter_reading_occurrence,
                    first_consumer_party_name, second_consumer_party_name,
                    meter_number, consumer_start_date, type_of_mp,
                    balance_supplier_name, postcode
                ) VALUES (
                    %(meteringPointId)s, %(streetCode)s, %(streetName)s, %(buildingNumber)s,
                    %(floorId)s, %(roomId)s, %(citySubDivisionName)s, %(municipalityCode)s,
                    %(locationDescription)s, %(settlementMethod)s, %(meterReadingOccurrence)s,
                    %(firstConsumerPartyName)s, %(secondConsumerPartyName)s,
                    %(meterNumber)s, %(consumerStartDate)s, %(typeOfMP)s,
                    %(balanceSupplierName)s, %(postcode)s
                )
                ON CONFLICT (metering_point_id) 
                DO UPDATE SET
                    street_code = EXCLUDED.street_code,
                    street_name = EXCLUDED.street_name,
                    building_number = EXCLUDED.building_number,
                    updated_at = CURRENT_TIMESTAMP;
                """

                cursor.execute(upsert_sql, metering_point_data)
                self.connection.commit()
                print(
                    f"Metering point {metering_point_data['meteringPointId']} stored successfully"
                )
                return True

        except Exception as e:
            print(f"Error storing metering point: {e}")
            self.connection.rollback()
            return False

    def store_consumption_readings(self, metering_point_id, consumption_data):
        """Store consumption readings with conflict handling"""
        if not self.connection:
            print("No database connection")
            return False

        try:
            with self.connection.cursor() as cursor:
                # Prepare batch insert with conflict handling
                insert_sql = """
                INSERT INTO consumption_readings (
                    metering_point_id, timestamp, consumption_kwh, quality, period_resolution
                ) VALUES (
                    %s, %s, %s, %s, %s
                )
                ON CONFLICT (metering_point_id, timestamp) 
                DO UPDATE SET
                    consumption_kwh = EXCLUDED.consumption_kwh,
                    quality = EXCLUDED.quality;
                """

                readings_data = []
                for reading in consumption_data:
                    readings_data.append(
                        (
                            metering_point_id,
                            reading["timestamp"],
                            reading["consumption_kwh"],
                            reading["quality"],
                            reading.get("period_resolution", "PT1H"),
                        )
                    )

                cursor.executemany(insert_sql, readings_data)
                self.connection.commit()
                print(f"Stored {len(readings_data)} consumption readings")
                return True

        except Exception as e:
            print(f"Error storing consumption readings: {e}")
            self.connection.rollback()
            return False

    def close(self):
        """Close database connection"""
        if self.connection:
            self.connection.close()
            print("Database connection closed")


# Example usage:
print("EnergyDataStorage class defined. Use it like this:")
print("""
# Initialize storage
storage = EnergyDataStorage(DB_CONFIG)
storage.connect()
storage.create_tables()

# Store metering point data
storage.store_metering_point(metering_point_data)

# Store consumption readings
storage.store_consumption_readings(metering_point_id, readings_list)

storage.close()
""")
In [ ]:
# Function to transform API data for database storage


def transform_consumption_data_for_db(consumption_json):
    """
    Transform the Eloverblik API response into database-ready format
    """
    if not consumption_json or "result" not in consumption_json:
        return None, []

    try:
        # Extract metering point data
        market_document = consumption_json["result"][0]["MyEnergyData_MarketDocument"]
        time_series = market_document["TimeSeries"][0]
        metering_point_id = time_series["mRID"]

        # For demonstration, we'll create a basic metering point record
        # In practice, you'd get this from the metering points API call
        metering_point_data = {
            "meteringPointId": metering_point_id,
            "streetCode": None,
            "streetName": None,
            "buildingNumber": None,
            "floorId": None,
            "roomId": None,
            "citySubDivisionName": None,
            "municipalityCode": None,
            "locationDescription": None,
            "settlementMethod": None,
            "meterReadingOccurrence": None,
            "firstConsumerPartyName": None,
            "secondConsumerPartyName": None,
            "meterNumber": None,
            "consumerStartDate": None,
            "typeOfMP": None,
            "balanceSupplierName": None,
            "postcode": None,
        }

        # Extract consumption readings
        consumption_readings = []
        periods = time_series["Period"]

        for period in periods:
            period_start = datetime.fromisoformat(
                period["timeInterval"]["start"].replace("Z", "+00:00")
            )
            resolution = period["resolution"]

            if "Point" in period:
                for point in period["Point"]:
                    # Calculate the actual timestamp for this point
                    position = int(point["position"])
                    # Position is 1-based, so subtract 1 to get hours offset
                    hours_offset = position - 1

                    point_timestamp = period_start + timedelta(hours=hours_offset)

                    reading = {
                        "timestamp": point_timestamp,
                        "consumption_kwh": float(point["out_Quantity.quantity"]),
                        "quality": point["out_Quantity.quality"],
                        "period_resolution": resolution,
                    }
                    consumption_readings.append(reading)

        return metering_point_data, consumption_readings

    except Exception as e:
        print(f"Error transforming data: {e}")
        return None, []


# Test the transformation with our existing data
if "consumption_json" in locals() and consumption_json:
    metering_point, readings = transform_consumption_data_for_db(consumption_json)

    if metering_point and readings:
        print(
            f"Transformed data for metering point: {metering_point['meteringPointId']}"
        )
        print(f"Number of readings: {len(readings)}")
        print("\nSample readings:")
        for i, reading in enumerate(readings[:3]):
            print(
                f"  {i + 1}. {reading['timestamp']}: {reading['consumption_kwh']} kWh (Quality: {reading['quality']})"
            )
        print("  ...")
    else:
        print("Failed to transform data")
else:
    print("No consumption data available to transform")