(migrate_till: int = 30)
| 5 | |
| 6 | |
| 7 | def migrate_feature_evaluations(migrate_till: int = 30) -> None: |
| 8 | query_api = InfluxDBWrapper.get_client().query_api() |
| 9 | read_bucket = InfluxDBWrapper.get_downsampled_bucket(DownsampleSize.FIFTEEN_MINUTES) |
| 10 | |
| 11 | for i in range(migrate_till): |
| 12 | range_start = f"-{i + 1}d" |
| 13 | range_stop = f"-{i}d" |
| 14 | query = ( |
| 15 | f'from (bucket: "{read_bucket}") ' |
| 16 | f"|> range(start: {range_start}, stop: {range_stop}) " |
| 17 | f'|> filter(fn: (r) => r._measurement == "feature_evaluation")' |
| 18 | ) |
| 19 | |
| 20 | result = query_api.query(query) |
| 21 | |
| 22 | feature_evaluations = [] |
| 23 | for table in result: |
| 24 | for record in table.records: |
| 25 | feature_evaluations.append( |
| 26 | FeatureEvaluationBucket( |
| 27 | feature_name=record.values["feature_id"], |
| 28 | bucket_size=ANALYTICS_READ_BUCKET_SIZE, |
| 29 | created_at=record.get_time(), |
| 30 | total_count=record.get_value(), |
| 31 | environment_id=record.values["environment_id"], |
| 32 | ) |
| 33 | ) |
| 34 | FeatureEvaluationBucket.objects.bulk_create(feature_evaluations) |
searching dependent graphs…