1
votes

Nettoyer la mémoire dans le workflow Drake R

J'ai créé un workflow de séries chronologiques massives (modèles 4273 * 10) pour 4273 séries chronologiques par semaine dans drake.

À l'origine, j'ai tenté de créer le flux de travail complet en utilisant le package fable. Ce qui est assez pratique pour entraîner des modèles pour des tsibbles groupés, mais après différents essais, j'ai eu beaucoup de problèmes avec la gestion de la mémoire. Mon serveur RStudio avec 32 cœurs et 244 Go de RAM plantait constamment, spécialement lorsque j'essayais de sérialiser les modèles.

À cause de cela, j'ai complètement craché mon flux de travail afin d'identifier les goulots d'étranglement provenant de:

 entrez la description de l'image ici

À:

 entrez la description de l'image ici

Puis à:

 entrez la description de l'image ici

Enfin à:

 entrez la description de l'image ici

Dans mon code d'entraînement (exemple prophet_multiplicative), j'utilise le futur package pour entraîner ces multiples modèles de fables, puis calculer la précision et les enregistrer. Cependant, je ne sais pas comment supprimer cet objet du workflow de drake par la suite:

  • Dois-je simplement supprimer l'objet à l'aide de rm?
  • Est-il possible dans drake d'avoir des environnements séparés pour chacun des composants du flux de travail?
  • Est-ce la bonne solution?

Mon idée est d'exécuter chacune des techniques individuelles de manière sérielle pendant que les 4273 modèles pour une technique spécifique sont entraînés en parallèle. Ce faisant, je m'attends à ne pas planter le serveur, puis une fois tous mes modèles formés, je peux rejoindre les métriques de précision, choisir le meilleur modèle pour chacune de mes séries chronologiques, puis découper chacun des fichiers binaires individuels pour pouvoir produire le prévisions.

Toutes les suggestions concernant mon approche sont les bienvenues. Veuillez noter que je suis assez limité en ressources matérielles, donc obtenir un serveur plus gros n'est pas une option.

BR / E


0 commentaires

3 Réponses :


2
votes

Il y a toujours un compromis entre la mémoire et la vitesse. Pour économiser de la mémoire, nous devons décharger certaines cibles de la session, ce qui nous oblige souvent à prendre le temps de les lire à partir du stockage plus tard. Le comportement par défaut de drake est de privilégier la vitesse. Donc, dans votre cas, je définirais memory_strategy = "autoclean" et garbage_collection = TRUE dans make () et les fonctions associées. Le manuel d'utilisation comporte un chapitre consacré à la gestion de la mémoire: https://books.ropensci.org/drake /memory.html .

De plus, je recommande de renvoyer de petites cibles lorsque cela est possible. Ainsi, au lieu de renvoyer un modèle ajusté entier, vous pouvez à la place renvoyer une petite trame de données de résumés de modèle, ce qui sera plus doux pour la mémoire et le stockage. En plus de cela, vous pouvez choisir l'un des formats de stockage spécialisés sur https://books.ropensci.org/drake/plans.html#special-data-formats-for-targets pour gagner encore plus d'efficacité.


0 commentaires

0
votes

garbage_collection = TRUE est déjà défini. J'essaierai d'ajouter l'autoclean. En ce qui concerne les formats de fichiers, j'enregistre mes modèles au format .qs avec la bibliothèque qs en utilisant la fonction save_model_x:

make(plan = plan, verbose = 2, 
     log_progress = TRUE,
     recover = TRUE,
     lock_envir = FALSE,
     garbage_collection = TRUE,
     memory_strategy = "autoclean")

Dans mon plan, cela est utilisé comme:

prophet_multiplicative = trainModels(input_data = processed_data, 
                               max_forecast_horizon = argument_parser$horizon,
                               max_multisession_cores = 6,
                               model_type = "prophet_multiplicative"),
  accuracy_prophet_multiplicative = accuracy_explorer(type = "train", models = prophet_multiplicative, 
                                                      max_forecast_horizon = argument_parser$horizon,
                                                      directory_out = "/data1/my_folder/"),
  saving_prophet_multiplicative = saveModels(models = prophet_multiplicative, 
                       directory_out = "/data1/my_folder/,
                       max_forecasting_horizon = argument_parser$horizon,
                       max_multisession_cores = 6)


6 commentaires

Faites-moi savoir si vous rencontrez toujours des problèmes de mémoire même après l'autoclean. Vous pouvez également regarder la stratégie de mémoire «aucun» si vous voulez vraiment prendre le contrôle manuel complet. L'un des avantages de drake est qu'il extrait les fichiers en tant qu'objets R et gère le stockage pour vous. Donc, si vous en avez envie, une alternative aux appels qsave () personnalisés est drake_plan (target (your_target, your_command (), format = "qs")) .


Là encore, si vous continuez à rencontrer des problèmes de mémoire, une tactique complètement opposée à la target (format = "qs") consiste à utiliser fichiers dynamiques pour tout. Avec les fichiers dynamiques, drake ne garde que le chemin du fichier en mémoire, pas l'objet lui-même, mais c'est à vous de lire manuellement l'objet en mémoire pour chaque cible qui l'utilise.


Salut Landau, pour une raison quelconque, ma dernière réponse n'a pas été publiée après votre commentaire. Maintenant, j'ai un nouveau problème, je pense que les modèles de train de fonctions ne s'évaluent pas correctement après l'ajout d'auto_clean.


La fonction se comporte-t-elle normalement en dehors de Drake? Si tel est le cas, publieriez-vous un exemple reproductible dans un autre fil de discussion. Cela semble être un problème dont je devrai me lancer pour résoudre le problème.


Salut Landau, Oui, la fonction fonctionne parfaitement en dehors de Drake. J'ajouterai l'exemple.


Souhaitez-vous ouvrir une toute nouvelle question Stack Overflow ou un problème GitHub? Je ne suis pas notifié lorsque vous publiez une nouvelle réponse / solution comme réponse sur cette page.



0
votes

Merci pour les réponses rapides. J'ai vraiment apprécié. Maintenant, je suis confronté à un autre problème, je laisse le script s'exécuter la nuit via nohup et j'ai trouvé ce qui suit dans les journaux:

trainModels <- function(input_data, max_forecast_horizon, model_type, max_multisession_cores) {

  options(future.globals.maxSize = 1500000000)
  future::plan(multisession, workers = max_multisession_cores) #breaking infrastructure once again ;)
  set.seed(666) # reproducibility
  
    if(max_forecast_horizon <= 104) {
      
      print(paste0("Training ", model_type, " models for forecasting horizon ", max_forecast_horizon))
      print(paste0("Using ", max_multisession_cores, " sessions from as future::plan()"))
      
      if(model_type == "prophet_multiplicative") {
        
        ts_models <- input_data %>% model(prophet = fable.prophet::prophet(snsr_val_clean ~ season("week", 2, type = "multiplicative") + 
                                                                             season("month", 2, type = "multiplicative")))
        
      } else if(model_type == "prophet_additive") {
        
        ts_models <- input_data %>% model(prophet = fable.prophet::prophet(snsr_val_clean ~ season("week", 2, type = "additive") + 
                                                                             season("month", 2, type = "additive")))
        
      } else if(model_type == "auto.arima") {
        
        ts_models <- input_data %>% model(auto_arima = ARIMA(snsr_val_clean))
        
      } else if(model_type == "arima_with_yearly_fourier_components") {
        
        ts_models <- input_data %>% model(auto_arima_yf = ARIMA(snsr_val_clean ~ fourier("year", K = 2)))
        
      } else if(model_type == "arima_with_monthly_fourier_components") {
        
        ts_models <- input_data %>% model(auto_arima_mf = ARIMA(snsr_val_clean ~ fourier("month", K=2)))
        
      } else if(model_type == "regression_with_arima_errors") {
        
        ts_models <- input_data %>% model(auto_arima_mf_reg = ARIMA(snsr_val_clean ~ month + year  + quarter + qday + yday + week))
        
      } else if(model_type == "tslm") {
    
        ts_models <- input_data %>% model(tslm_reg_all = TSLM(snsr_val_clean ~ year  + quarter + month + day + qday + yday + week + trend()))
     
      } else if(model_type == "theta") {
        
        ts_models <- input_data %>% model(theta = THETA(snsr_val_clean ~ season()))
        
      } else if(model_type == "ensemble") {
        
        ts_models <- input_data %>% model(ensemble =  combination_model(ARIMA(snsr_val_clean), 
                                              ARIMA(snsr_val_clean ~ fourier("month", K=2)),
                                              fable.prophet::prophet(snsr_val_clean ~ season("week", 2, type = "multiplicative") +
                                              season("month", 2, type = "multiplicative"), 
                                              theta = THETA(snsr_val_clean ~ season()), 
                                              tslm_reg_all = TSLM(snsr_val_clean ~ year  + quarter + month + day + qday + yday + week + trend())))
            )
        
      }
      
    } 
  
    else if(max_forecast_horizon > 104) {
      
        print(paste0("Training ", model_type, " models for forecasting horizon ", max_forecast_horizon))
        print(paste0("Using ", max_multisession_cores, " sessions from as future::plan()"))
        
        
        if(model_type == "prophet_multiplicative") {
          
          ts_models <- input_data %>% model(prophet = fable.prophet::prophet(snsr_val_clean ~ season("month", 2, type = "multiplicative") + 
                                                                               season("month", 2, type = "multiplicative")))
          
        } else if(model_type == "prophet_additive") {
          
          ts_models <- input_data %>% model(prophet = fable.prophet::prophet(snsr_val_clean ~ season("month", 2, type = "additive") + 
                                                                               season("year", 2, type = "additive")))
          
        } else if(model_type == "auto.arima") {
          
          ts_models <- input_data %>% model(auto_arima = ARIMA(snsr_val_clean))
          
        } else if(model_type == "arima_with_yearly_fourier_components") {
          
          ts_models <- input_data %>% model(auto_arima_yf = ARIMA(snsr_val_clean ~ fourier("year", K = 2)))
          
        } else if(model_type == "arima_with_monthly_fourier_components") {
          
          ts_models <- input_data %>% model(auto_arima_mf = ARIMA(snsr_val_clean ~ fourier("month", K=2)))
          
        } else if(model_type == "regression_with_arima_errors") {
          
          ts_models <- input_data %>% model(auto_arima_mf_reg = ARIMA(snsr_val_clean ~ month + year  + quarter + qday + yday))
          
        } else if(model_type == "tslm") {
          
          ts_models <- input_data %>% model(tslm_reg_all = TSLM(snsr_val_clean ~ year  + quarter + month + day + qday + yday + trend()))
          
        } else if(model_type == "theta") {
          
          ts_models <- input_data %>% model(theta = THETA(snsr_val_clean ~ season()))
          
        } else if(model_type == "ensemble") {
          
          ts_models <- input_data %>% model(ensemble =  combination_model(ARIMA(snsr_val_clean), 
                                                                          ARIMA(snsr_val_clean ~ fourier("month", K=2)),
                                                                          fable.prophet::prophet(snsr_val_clean ~ season("month", 2, type = "multiplicative") +
                                                                          season("year", 2, type = "multiplicative"),
                                                                          theta = THETA(snsr_val_clean ~ season()), 
                                                                          tslm_reg_all = TSLM(snsr_val_clean ~ year  + quarter + month + day + qday + 
                                                                                                yday  + trend())))
          )
          
        }
    }
  
  return(ts_models)
}

L'objet ts_models est l'objet en cours de création dans mes scripts d'entraînement et il est essentiellement ce que ma fonction trainModels renvoie. Il me semble que le paramètre de données d'entrée est peut-être propre et que c'est la raison pour laquelle il échoue?

Une autre question pour une raison quelconque, mon modèle n'est pas enregistré après l'entraînement des modèles thetha. Existe-t-il un moyen de spécifier que drake ne passe pas au modèle suivant tant qu'il n'a pas calculé la précision de l'un et enregistré le fichier .qs?

Ma fonction d'entraînement est la suivante:

[1] "DB PROD Connected"
[1] "DB PROD Connected"
[1] "Getting RAW data"
[1] "Maximum forecasting horizon is 52, fetching weekly data"
[1] "Removing duplicates if we have them"
[1] "Original data has 1860590 rows"
[1] "Data without duplicates has 1837995 rows"
`summarise()` regrouping output by 'A', 'B' (override with `.groups` argument)
[1] "Removing non active customers"
[1] "Data without duplicates and without active customers has 1654483 rows"
0.398 sec elapsed
[1] "Removing customers with last data older than 1.5 years"
[1] "Data without duplicates, customers that are not active and old customers has 1268610 rows"
0.223 sec elapsed
[1] "Augmenting data"
12.103 sec elapsed
[1] "Creating tsibble"
7.185 sec elapsed
[1] "Filling gaps for not breaking groups"
9.568 sec elapsed
[1] "Training theta models for forecasting horizon 52"
[1] "Using 12 sessions from as future::plan()"
Repacking large object
[1] "Training auto_arima models for forecasting horizon 52"
[1] "Using 12 sessions from as future::plan()"
Error: target auto_arima failed.
diagnose(auto_arima)error$message:
  object 'ts_models' not found
diagnose(auto_arima)error$calls:
  1. └─global::trainModels(...)
In addition: Warning message:
9 errors (2 unique) encountered for theta
[3] function cannot be evaluated at initial parameters
[6] Not enough data to estimate this ETS model.

Execution halted
            

BR / E


4 commentaires

Je vois ça maintenant. Pour résoudre le problème, il serait vraiment utile de voir une version réduite de l'ensemble du projet qui reproduit l'erreur. La fonction semble bonne en un coup d'œil, mais j'ai également besoin de voir le plan et d'autres codes contextuels et être capable de tout exécuter moi-même.


Bonjour Landau: github.com/ropensci/drake/issues/1293 , plan en détail là


Merci, je vais jeter un oeil.


Merci à vous, un support vraiment génial, très heureux d'utiliser drake!