MongoDB
Ingest Meters from MongoDB
Amberflo supports metering data ingestion from MongoDB, allowing you to set up metering alongside your current system without disruption.
If your system already uses MongoDB, you can configure Amberflo to ingest data without impacting existing functionality.
Recommended Methods for MongoDB Integration
We recommend the following approaches to ingest data from MongoDB into Amberflo:
1.Using Logstash
Logstash is an open-source data collection engine with real-time data pipelining capabilities.
- It can unify data from different sources, normalize it, and forward it to your target systems.
- By using Logstash’s input, output, filtering, and buffering capabilities, you can reliably move data from MongoDB to Amberflo.
Plugin Setup:
- Use either the MongoDB input plugin or the JDBC input plugin
- Configure the S3 output plugin to write to the Amberflo-provided S3 bucket
Once data is in the S3 bucket, Amberflo will automatically pick it up for ingestion. Elastic Logstash https://www.elastic.co/guide/en/logstash/current/plugins-inputs-jdbc.html https://github.com/phutchins/logstash-input-mongodb
2.Query MongoDB using code
You can also query MongoDB programmatically using its SDK and then use the Amberflo SDK to ingest the retrieved data.
Example (Python)
- Use PyMongo to query data from MongoDB.
- Use the Amberflo Python SDK to ingest the results as meter records.
This approach gives you full control over:
- Filtering and transforming data before ingestion
- Handling new data only, by persisting the last processed timestamp or ID
Amberflo automatically performs deduplication, so repeated ingestion of the same data will not create duplicates.
3.Use MongoDB's Change data stream
With the Amberflo SDK, you can use the MongoDB real-time data stream to act upon every insert: https://developer.mongodb.com/quickstart/python-change-streams/ Python SDK
Set up virtual environment and install SDK.# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up customer API client and create a customer
from metering.customer import CustomerApiClient, create_customer_payload
customer_api = CustomerApiClient(os.environ.get("API_KEY"))# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up customer API client and create a customer
from metering.customer import CustomerApiClient, create_customer_payload
customer_api = CustomerApiClient(os.environ.get("API_KEY"))
message = create_customer_payload(
customer_id="customer-123",
customer_name="Customer 123",
# email is optional
customer_email="[email protected]",
# traits are optional. They can be used as filters or aggregation buckets
traits={
"region": "us-east-1",
},
)
customer = customer_api.add_or_update(message)# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up customer API client and create a customer
from metering.customer import CustomerApiClient, create_customer_payload
customer_api = CustomerApiClient(os.environ.get("API_KEY"))
message = create_customer_payload(
customer_id="customer-123",
customer_name="Customer 123",
# email is optional
customer_email="[email protected]",
# traits are optional. They can be used as filters or aggregation buckets
traits={
"region": "us-east-1",
},
)
customer = customer_api.add_or_update(message)
# Set up ingest API client and ingest a meter
from time import time
from metering.ingest import create_ingest_client
ingest_api = create_ingest_client(api_key=API_KEY)# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up customer API client and create a customer
from metering.customer import CustomerApiClient, create_customer_payload
customer_api = CustomerApiClient(os.environ.get("API_KEY"))
message = create_customer_payload(
customer_id="customer-123",
customer_name="Customer 123",
# email is optional
customer_email="[email protected]",
# traits are optional. They can be used as filters or aggregation buckets
traits={
"region": "us-east-1",
},
)
customer = customer_api.add_or_update(message)
# Set up ingest API client and ingest a meter
from time import time
from metering.ingest import create_ingest_client
ingest_api = create_ingest_client(api_key=API_KEY)
ingest_api.meter(
meter_api_name="api-calls",
meter_value=1,
meter_time_in_millis=int(time() * 1000),
customer_id=customer["customerId"],
)# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up customer API client and create a customer
from metering.customer import CustomerApiClient, create_customer_payload
customer_api = CustomerApiClient(os.environ.get("API_KEY"))
message = create_customer_payload(
customer_id="customer-123",
customer_name="Customer 123",
# email is optional
customer_email="[email protected]",
# traits are optional. They can be used as filters or aggregation buckets
traits={
"region": "us-east-1",
},
)
customer = customer_api.add_or_update(message)
# Set up ingest API client and ingest a meter
from time import time
from metering.ingest import create_ingest_client
ingest_api = create_ingest_client(api_key=API_KEY)
ingest_api.meter(
meter_api_name="api-calls",
meter_value=1,
meter_time_in_millis=int(time() * 1000),
customer_id=customer["customerId"],
)
# Shut down ingest API client
ingest_api.shutdown()# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up customer API client and create a customer
from metering.customer import CustomerApiClient, create_customer_payload
customer_api = CustomerApiClient(os.environ.get("API_KEY"))
message = create_customer_payload(
customer_id="customer-123",
customer_name="Customer 123",
# email is optional
customer_email="[email protected]",
# traits are optional. They can be used as filters or aggregation buckets
traits={
"region": "us-east-1",
},
)
customer = customer_api.add_or_update(message)
# Set up ingest API client and ingest a meter
from time import time
from metering.ingest import create_ingest_client
ingest_api = create_ingest_client(api_key=API_KEY)
ingest_api.meter(
meter_api_name="api-calls",
meter_value=1,
meter_time_in_millis=int(time() * 1000),
customer_id=customer["customerId"],
)
# Shut down ingest API client
ingest_api.shutdown()
# Setup usage API client and get a usage report
from metering.usage import (AggregationType, Take, TimeGroupingInterval,
TimeRange, UsageApiClient, create_usage_query)
usage_api = UsageApiClient(os.environ.get("API_KEY"))# Set up virtual environment and install SDK
#
# $ python3 -m virtualenv venv
# $ . venv/bin/activate
# $ pip install amberflo-metering-python
import os
# Sign up for free at https://ui.amberflo.io/ to get an API key
API_KEY = os.environ["AMBERFLO_API_KEY"]
# Set up customer API client and create a customer
from metering.customer import CustomerApiClient, create_customer_payload
customer_api = CustomerApiClient(os.environ.get("API_KEY"))
message = create_customer_payload(
customer_id="customer-123",
customer_name="Customer 123",
# email is optional
customer_email="[email protected]",
# traits are optional. They can be used as filters or aggregation buckets
traits={
"region": "us-east-1",
},
)
customer = customer_api.add_or_update(message)
# Set up ingest API client and ingest a meter
from time import time
from metering.ingest import create_ingest_client
ingest_api = create_ingest_client(api_key=API_KEY)
ingest_api.meter(
meter_api_name="api-calls",
meter_value=1,
meter_time_in_millis=int(time() * 1000),
customer_id=customer["customerId"],
)
# Shut down ingest API client
ingest_api.shutdown()
# Setup usage API client and get a usage report
from metering.usage import (AggregationType, Take, TimeGroupingInterval,
TimeRange, UsageApiClient, create_usage_query)
usage_api = UsageApiClient(os.environ.get("API_KEY"))
since_two_days_ago = TimeRange(int(time()) - 60 * 60 * 24 * 2)
query = create_usage_query(
meter_api_name="api-calls",
aggregation=AggregationType.SUM,
time_grouping_interval=TimeGroupingInterval.DAY,
time_range=since_two_days_ago,
group_by=["customerId"],
usage_filter={"customerId": ["customer-123"]},
take=Take(limit=10, is_ascending=False),
)
report = usage_api.get(query)