Skip to content

Commit 0ed4b9b

Browse files
authored
Merge pull request #604 from ddps-lab/azure-collector
Add Azure T2/T3 Logic
2 parents b5403bd + b0d9133 commit 0ed4b9b

4 files changed

Lines changed: 82 additions & 8 deletions

File tree

collector/spot-dataset/azure/lambda/current_collector/lambda_function_sps.py

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
from sps_module import sps_shared_resources
77
from utils.merge_df import merge_if_saving_price_sps_df
88
from utils.upload_data import update_latest, save_raw, upload_timestream, query_selector, upload_cloudwatch
9-
from utils.compare_data import compare_sps
9+
from utils.compare_data import compare_sps, compare_max_instance
1010
from utils.pub_service import send_slack_message, Logger, S3, AZURE_CONST
1111
from utils.azure_auth import get_sps_token_and_subscriptions
1212

@@ -164,13 +164,17 @@ def handle_res_df_for_spotlake(price_saving_if_df, sps_df, time_datetime, desire
164164
f"{AZURE_CONST.S3_LATEST_ALL_DATA_AVAILABILITY_ZONE_TRUE_PKL_GZIP_SAVE_PATH}", 'pkl.gz')
165165

166166
workload_cols = ['InstanceTier', 'InstanceType', 'Region', 'AvailabilityZone', 'DesiredCount']
167-
feature_cols = ['OndemandPrice', 'SpotPrice', 'IF', 'Score', 'SPS_Update_Time']
167+
feature_cols = ['OndemandPrice', 'SpotPrice', 'IF', 'Score', 'SPS_Update_Time', 'T2', 'T3']
168168

169169
query_success = timestream_success = cloudwatch_success = \
170170
update_latest_success = save_raw_success = False
171171

172172
if prev_availability_zone_true_all_data_df is not None and not prev_availability_zone_true_all_data_df.empty:
173-
prev_availability_zone_true_all_data_df.drop(columns=['id'], inplace=True)
173+
prev_availability_zone_true_all_data_df.drop(columns=['id'], inplace=True, errors='ignore')
174+
175+
# Apply T2/T3 Aggregation Logic
176+
sps_merged_df = compare_max_instance(prev_availability_zone_true_all_data_df, sps_merged_df, current_desired_count)
177+
174178
changed_df = compare_sps(prev_availability_zone_true_all_data_df, sps_merged_df, workload_cols, feature_cols)
175179

176180
query_success = query_selector(changed_df)

collector/spot-dataset/azure/lambda/current_collector/load_sps.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -213,7 +213,9 @@ def execute_spot_placement_score_task_by_parameter_pool_df(api_calls_df, desired
213213
score.get("sku", ""), {}).get("InstanceTier"),
214214
"InstanceType": SS_Resources.region_map_and_instance_map_tmp['instance_map'].get(
215215
score.get("sku", ""), {}).get("InstanceTypeOld"),
216-
"Score": score.get("score")
216+
"Score": score.get("score"),
217+
"T3": desired_count if score.get("score") == 3 else 0,
218+
"T2": desired_count if score.get("score") == 2 else 0
217219
}
218220
if availability_zones is True:
219221
score_data["AvailabilityZone"] = score.get("availabilityZone", "Single")

collector/spot-dataset/azure/lambda/current_collector/utils/compare_data.py

Lines changed: 65 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,4 +50,68 @@ def compare_sps(previous_df, current_df, workload_cols, feature_cols):
5050

5151
changed_df = changed_df[current_df.columns]
5252

53-
return changed_df if not changed_df.empty else None
53+
return changed_df if not changed_df.empty else None
54+
55+
def compare_max_instance(previous_df, new_df, target_capacity):
56+
fallback_dict = {50:45, 45:40, 40:35, 35:30, 30:25, 25:20, 20:15, 15:10, 10:5, 5:1, 1:0}
57+
fallback_val = fallback_dict.get(target_capacity, 0)
58+
59+
merged_df = pd.merge(
60+
new_df,
61+
previous_df[["InstanceType", "Region", "AvailabilityZone", "DesiredCount", "Score", "T3", "T2"]],
62+
on=["InstanceType", "Region", "AvailabilityZone", "DesiredCount"],
63+
how="left",
64+
suffixes=("", "_prev")
65+
)
66+
67+
# Fix SPS when single node SPS
68+
if target_capacity == 1:
69+
merged_df["Score"] = merged_df["Score"].combine_first(merged_df["Score_prev"])
70+
71+
# Merge single node SPS with multi node SPS if (multi node SPS) > (single node SPS)
72+
# Note: Score strings "3", "2", "1" are comparable.
73+
# But need to handle N/A or types. Assuming Score is int or convertible.
74+
# Azure Score is int from load_sps.
75+
# previous_df Score might be string if read from file? feature_cols convert to str in compare_sps but here we read raw df.
76+
77+
# Ensure Score types are compatible (float/int)
78+
merged_df["Score"] = pd.to_numeric(merged_df["Score"], errors='coerce')
79+
merged_df["Score_prev"] = pd.to_numeric(merged_df["Score_prev"], errors='coerce')
80+
merged_df["T3"] = pd.to_numeric(merged_df["T3"], errors='coerce').fillna(0)
81+
merged_df["T3_prev"] = pd.to_numeric(merged_df["T3_prev"], errors='coerce').fillna(0)
82+
merged_df["T2"] = pd.to_numeric(merged_df["T2"], errors='coerce').fillna(0)
83+
merged_df["T2_prev"] = pd.to_numeric(merged_df["T2_prev"], errors='coerce').fillna(0)
84+
85+
merged_df.loc[(merged_df["Score"] > merged_df["Score_prev"]), "Score_prev"] = merged_df["Score"]
86+
87+
# Calculate T3
88+
# Use numpy where.
89+
merged_df["T3"] = np.where(
90+
merged_df["Score"] >= 3,
91+
np.maximum(merged_df["T3"], merged_df["T3_prev"]),
92+
np.minimum(fallback_val, merged_df["T3_prev"])
93+
)
94+
95+
# Calculate T2
96+
merged_df["T2"] = np.where(
97+
merged_df["Score"] >= 2,
98+
np.maximum(merged_df["T2"], merged_df["T2_prev"]),
99+
np.minimum(fallback_val, merged_df["T2_prev"])
100+
)
101+
102+
if target_capacity == 1:
103+
merged_df.loc[merged_df["Score"] <= 2, "T3"] = 0
104+
merged_df.loc[merged_df["Score"] < 2, "T2"] = 0
105+
else:
106+
merged_df.loc[merged_df["Score_prev"] <= 2, "T3"] = 0
107+
merged_df.loc[merged_df["Score_prev"] < 2, "T2"] = 0
108+
# Fix SPS to Single node SPS ? mimic AWS "merged_df["SPS"] = merged_df["SPS_prev"]"
109+
# AWS comment: "Fix SPS to Single node SPS" - this logic seems specific to assuming single node fallback.
110+
# But let's copy logic:
111+
merged_df["Score"] = merged_df["Score_prev"]
112+
113+
# Convert to standard types if needed
114+
# Drop unnecessary columns
115+
merged_df.drop(columns=["T3_prev", "T2_prev", "Score_prev"], errors='ignore', inplace=True)
116+
117+
return merged_df

collector/spot-dataset/azure/lambda/current_collector/utils/upload_data.py

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -137,7 +137,7 @@ def upload_timestream(data, time_datetime):
137137
try:
138138
data = data.copy()
139139
data = data[["InstanceTier", "InstanceType", "Region", "OndemandPrice", "SpotPrice", "Savings", "IF",
140-
"DesiredCount", "AvailabilityZone", "Score", "SPS_Update_Time"]]
140+
"DesiredCount", "AvailabilityZone", "Score", "SPS_Update_Time", "T2", "T3"]]
141141

142142
fill_values = {
143143
"InstanceTier": 'N/A',
@@ -150,7 +150,9 @@ def upload_timestream(data, time_datetime):
150150
'DesiredCount': -1,
151151
'AvailabilityZone': 'N/A',
152152
'Score': 'N/A',
153-
'SPS_Update_Time': 'N/A'
153+
'SPS_Update_Time': 'N/A',
154+
'T2': 0,
155+
'T3': 0
154156
}
155157
data = data.fillna(fill_values)
156158

@@ -177,7 +179,9 @@ def upload_timestream(data, time_datetime):
177179
('SpotPrice', 'DOUBLE'),
178180
('IF', 'DOUBLE'),
179181
('Score', 'VARCHAR'),
180-
('SPS_Update_Time', 'VARCHAR')
182+
('SPS_Update_Time', 'VARCHAR'),
183+
('T2', 'DOUBLE'),
184+
('T3', 'DOUBLE')
181185
]
182186

183187
for column, value_type in measure_columns:

0 commit comments

Comments
 (0)