Multi-key GroupBy с общими данными на одном ключеPython

Программы на Python
Anonymous
Multi-key GroupBy с общими данными на одном ключе

Сообщение Anonymous »

Я работаю с большим набором данных, который включает в себя несколько уникальных групп данных, идентифицированных по дате и идентификатору группы. Каждая группа содержит несколько идентификаторов, каждый из которых имеет несколько атрибутов. Вот упрощенная структура моих данных:

Код: Выделить всё

| date       | group_id | inner_id | attr_a | attr_b | attr_c |
|------------|----------|----------|--------|--------|--------|
| 2023-06-01 | A1       | 001      | val    | val    | val    |
| 2023-06-01 | A1       | 002      | val    | val    | val    |
...
Кроме того, для каждой даты у меня есть большая связанная с ней матрица:

Код: Выделить всё

| date       | matrix       |
|------------|--------------|
| 2023-06-01 | [[...], ...] |
...
Мне нужно применить функцию для каждой даты и group_id, которая обрабатывает данные, используя как атрибуты группы, так и матрицу, связанную с этой датой. Функция выглядит следующим образом:

Код: Выделить всё

def run(group_data: pd.DataFrame, matrix) -> pd.DataFrame:
# process data
return processed_data
Здесь group_data содержит атрибуты для конкретной группы:

Код: Выделить всё

| inner_id | attr_a | attr_b | attr_c |
|----------|--------|--------|--------|
| 001      | val    | val    | val    |
...
Вот моя текущая реализация, она работает, но я могу запускать только ~200 дат одновременно, потому что я передаю все данные всем работникам (у меня ~2 тыс. дат, ~100 групп на date, ~150 внутренних элементов на группу)

Код: Выделить всё

def calculate_metrics(data: DataFrame, matrices: DataFrame) -> DataFrame:
# Convert matrices to a dictionary mapping dates to matrix
date_matrices = matrices.rdd.collectAsMap()

# Broadcast the matrices
broadcasted_matrices = spark_context.broadcast(date_matrices)

# Function to apply calculations
def apply_calculation(group_key: Tuple[str, str], data_group: pd.DataFrame) -> pd.DataFrame:
date = group_key[1]
return custom_calculation_function(broadcasted_matrices.value[date], data_group)

# Apply the function to each group
return data.groupby('group_id', 'date').applyInPandas(apply_calculation, schema_of_result)
Как оптимизировать эти вычисления для эффективного распараллеливания обработки и обеспечения того, чтобы матрицы не загружались в память избыточно больше, чем необходимо?


Подробнее здесь: https://stackoverflow.com/questions/786 ... on-one-key

Вернуться в «Python»