feat: bigquery_fetcher entrypoint - #67
Conversation
|
|
||
| [project.optional-dependencies] | ||
| # Only needed for the `bigquery_fetcher` entrypoint | ||
| bigquery = ['google-cloud-bigquery>=3.0.0'] |
There was a problem hiding this comment.
Do we need to update the uv.lock as well? Although not sure if that's used anywhere in practice, maybe only for development.
| bucket_time = now if timestamp is None else timestamp | ||
| key = UsageKey( | ||
| floor(now / self.__granularity_sec) * self.__granularity_sec, | ||
| floor(bucket_time / self.__granularity_sec) | ||
| * self.__granularity_sec, |
There was a problem hiding this comment.
Maybe better?
| bucket_time = now if timestamp is None else timestamp | |
| key = UsageKey( | |
| floor(now / self.__granularity_sec) * self.__granularity_sec, | |
| floor(bucket_time / self.__granularity_sec) | |
| * self.__granularity_sec, | |
| bucket_time = timestamp if timestamp is not None else floor(now / self.__granularity_sec) * self.__granularity_sec | |
| key = UsageKey( | |
| bucket_time, |
There was a problem hiding this comment.
this only applies the bucketing to the now timestamp, not the passed-in one. but i agree the multi-line math expression in the arg list is ugly. i have cleaned it up
| query_list = parse_and_assert_query_file(query_file) | ||
| record_list: List[UsageAccumulatorRecord] = [] | ||
| for query_dict in query_list: | ||
| query = query_dict["query"] |
There was a problem hiding this comment.
I hope the query is pretty concise, otherwise I don't really love the idea of putting a large sql query in a JSON.
There was a problem hiding this comment.
the Bigtable query for us will look something like:
SELECT
SPLIT(rowkey, '/')[OFFSET(0)] AS app_feature, -- TODO mapping
SUM(
-- `m` field size: object metadata
IFNULL(BYTE_LENGTH(COALESCE(fg.m.cell.value, fm.m.cell.value)), 0)
-- `p` field size: object payload
-- check `m` for precomputed size, fallback to scanning `p`
+ IFNULL(
COALESCE(
LAX_INT64(
JSON_QUERY(
TO_JSON(COALESCE(fg.m.cell.value, fm.m.cell.value)), "$.size")),
BYTE_LENGTH(COALESCE(fg.p.cell.value, fm.p.cell.value))),
0)
-- `t` field size: tombstone metadata
+ IFNULL(BYTE_LENGTH(COALESCE(fg.t.cell.value, fm.t.cell.value)), 0)
-- `r` field size: tombstone
+ IFNULL(BYTE_LENGTH(COALESCE(fg.r.cell.value, fm.r.cell.value)), 0))
AS amount
FROM `eng-dev-sbx--files-1.objectstore.bigtable_objectstore`
GROUP BY usecase
ORDER BY total_bytes DESC;
GCS one should be shorter bc size is precomputed
lcian
left a comment
There was a problem hiding this comment.
Left some suggestions. Please address the CI checks as well.
a160ba5 to
41bdd31
Compare
41bdd31 to
7cb9809
Compare
|
|
||
| bq_client = bigquery.Client() | ||
|
|
||
| main( | ||
| args.query_file, | ||
| UsageAccumulator(kafka_config=kafka_config), | ||
| bq_client, |
There was a problem hiding this comment.
Bug: When dry_run=True, the usage_accumulator is not closed, leading to a resource leak from unclosed Kafka producer connections.
Severity: MEDIUM
Suggested Fix
Call usage_accumulator.close() within the if dry_run: block in the main function, after logging the records. This ensures the Kafka producer's resources are always cleaned up, mirroring the behavior of the non-dry-run path.
Prompt for AI Agent
Review the code at the location below. A potential bug has been identified by an AI
agent. Verify if this is a real issue. If it is, propose a fix; if not, explain why it's
not valid.
Location: py/usageaccountant/bigquery_fetcher.py#L179-L185
Potential issue: When the `main` function is executed with `dry_run=True`, the
`usage_accumulator` is instantiated but its `close()` method is never called. This
results in a resource leak, as the underlying `KafkaProducer` maintains open connections
to Kafka brokers that are not released. For a tool that might be run repeatedly in CI/CD
pipelines for validation, these unclosed connections could accumulate and lead to
resource exhaustion.
Also affects:
py/usageaccountant/datadog_fetcher.py:344~350
as title
timestampaccumulator for times when you're querying a snapshot. you can put in the snapshot timedatadog_fetcher.pyto newfetcher_utils.py(and same for tests)bigquery_fetcher.pyentrypoint which works basically the same as the datadog one but with a BQ clientDockerfileboilerplatei'm wiring some inventory data into BigQuery for objectstore COGS. reusing
usage-accountantto push BQ rollups through Kafka saves us a good deal of trouble.Ref FS-210