24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297 | class LuxorASAPManager:
@logger.catch
def __init__(self, reset=False, is_develop_mode=False):
self.today = dt.date.today()
self.__initialize(reset, is_develop_mode)
self.source_bases_exec_count = 0
def __reset(self):
self._register_assets()
self.save_funds()
def __initialize(self, reset, is_develop_mode):
self.is_develop_mode = is_develop_mode
# Criando db manager, inicia carregando todos os dados
self.lq = LuxorQuery(is_develop_mode=self.is_develop_mode)
self.dl = DataLoader(Path().absolute()/"LuxorDB"/"tables")
logger.info("LuxorDB: Dados Inicializados.")
self.dl.add_file_tracker(Path().absolute()/"source_bases"/"historico_trades_ativos.xlsx",
sheet_names={"Boletas": "trades", "Ativos": "assets", "VCs_prices":"hist_vc_px_last",
"asset_field_map": "asset_field_map", "field_map" : "field_map"}, index=False,
normalize_columns=True)
self.dl.load_file_if_modified(Path().absolute()/"source_bases"/"historico_trades_ativos.xlsx", export_to_blob=True)
# Carregando/criando fundos
self.funds = {}
if reset:
self.__reset()
else:
self.funds = self._load_funds()
def get_unique_funds(self):
""" Retorna uma lista de ocorrencia unicas de fundos na base de boletas """
unq_funds = list(self.lq.get_table("trades")["Fund"].unique())
unq_funds.remove("lipizzaner(eduardo)")
return unq_funds
def get_unique_assets(self, fund_filter=None):
""" Retorna lista de ocorrencias unicas de ativos. Se 'fund_filter' for informado, retorna as ocorrencias unicas nesse fundo."""
trades = self.lq.get_table("trades").sort_values(by="Date", kind="stable")
trades = self.lq.text_to_lowercase(trades)
trades["FRGN_KEY"] = trades["Asset"] + "_" + trades["Ticker"]
if fund_filter is None:
return trades["FRGN_KEY"].unique()
return trades.loc[ trades["Fund"] == fund_filter.lower(), "FRGN_KEY"].unique()
def get_asset_properties(self, asset_key):
""" Encontra as propriedades do ativo dado por 'asset_key'. Se nao for um ativo valido, retorna 'None'."""
valid_assets = self.lq.get_table("assets")
properties = valid_assets.loc[ valid_assets["Key"] == asset_key ]
if properties.empty:
return properties
return properties.iloc[0]
def _load_funds(self):
# Seleciona o arquivo mais recente e carrega os dados dos fundos
files = os.listdir("funds_data")
files.sort(reverse=True)
savefile = os.path.join("funds_data", files[0])
with open(savefile, "rb") as f:
funds = pickle.load(f)
return funds
def save_funds(self):
# salva os dados de cada fundo num pickle
attempts = 5
file_name = os.path.join("funds_data", "FundsData_"+str(self.today)+".pickle")
while attempts > 0:
try:
with open(file_name, "wb") as f:
pickle.dump(self.funds, f, pickle.HIGHEST_PROTOCOL)
attempts = 0
except PermissionError:
attempts -=1
logger.error(f"Erro ao tentar salvar FundsData_{str(self.today)}.pickle\nSerão feitas 5 tentativas...")
time.sleep(5)
def redirect_trades(self):
# Filtra o df de trades e redireciona para processamento individual em cada fundo.
trades = self.lq.get_table("trades")
for fund_name, fund in self.funds.items():
trades_df = trades.loc[ trades["Fund"] == fund_name ]
fund.set_trades_to_process(trades_df)
# TODO: Nao precisa usar o db_manager. Pode simplesmente salvar dentro da luxorDB.
def export_data_to_db(self, date_reference=None):
"""
Insere as tabelas processadas (posicoes, posicoes_dia) para luxor_db.
"""
if date_reference is None: date_reference=self.today
last_positions = []
hist_positions = []
cash_movements = [] # vamos agrupar as movimentacoes de caixa aqui para exportar
positions_by_bank = [] # vamos agrupar as posicoes por banco aqui para exportar
last_positions_by_bank = [] # vamos agrupar as posicoes por banco aqui para exportar
# Recuperamos a base de posicoes.. Mais recentes e historicas
for fund_name, fund_data in self.funds.items():
last_positions.append(fund_data.get_positions(all=False, date_reference=date_reference))
hist_positions.append(fund_data.get_positions(all=True))
cash_movements.append(fund_data.get_cash_movements())
positions_by_bank.append(fund_data.get_positions_by_bank())
last_positions_by_bank.append(fund_data.get_positions_by_bank(all=False, date_reference=date_reference))
last_update = dt.datetime.now().timestamp() # forcar atualizacao
# Concatenamos tudo numa tabela fato so
last_positions = pd.concat(last_positions)
# Exportamos as tabelas para luxorDB
# sobrescrevendo as movimentacoes na base (modificacoes ou nao, vamos sobrescrever pois base pode estar desatualizada)
logger.info(f"Exportando dados para a base de dados.")
self.dl.load_table_if_modified("last_positions", last_positions, last_update=last_update, normalize_columns=True,
export_to_blob=True)
#O mesmo eh feito com o restante
hist_positions = pd.concat(hist_positions)
self.dl.load_table_if_modified("hist_positions", hist_positions, last_update=last_update, normalize_columns=True,
export_to_blob=True)
# Exportamos as tabelas para luxorDB
cash_movements = pd.concat(cash_movements)
self.dl.load_table_if_modified("cash_movements", cash_movements, last_update=last_update, normalize_columns=True,
export_to_blob=True)
positions_by_bank = pd.concat(positions_by_bank)
self.dl.load_table_if_modified("hist_positions_by_bank", positions_by_bank, last_update=last_update, normalize_columns=True,
export_to_blob=True)
last_positions_by_bank = pd.concat(last_positions_by_bank)
self.dl.load_table_if_modified("last_positions_by_bank", last_positions_by_bank, last_update=last_update, normalize_columns=True,
export_to_blob=True)
def _register_assets(self):
# TODO PRECISAMOS SEPARAR POR CONTA BANCÁRIA!
# Restagando nomes unicos dos fundos
fund_names = self.get_unique_funds()
for name in fund_names:
self.funds[name] = Fund(name)
# Resgatando ativos unicos de cada fundo
fund_assets = self.get_unique_assets(fund_filter=name)
# Inicializamos cada ativo valido com suas respectivas props
for asset_key in fund_assets:
asset_props = self.get_asset_properties(asset_key)
if asset_props.empty:
continue # Ativo nao eh valido
props = {
"asset_type" : asset_props["Type"], "group" : asset_props["Group"], "ticker_bbg" : asset_props["Ticker_BBG"],
"curncy_exp" : asset_props["Currency Exposure"], "location" : asset_props["Location"]
}
self.funds[name].register_asset(asset_props["Asset"], asset_props["Ticker"], props, self.lq)
def _register_new_assets(self, fund_name):
assets_list = self.get_unique_assets(fund_filter=fund_name)
assets_already_registered = self.funds[fund_name].get_assets()
for asset_key in assets_list:
if asset_key not in assets_already_registered:
# Entao eh um ativo novo que precisamos registrar no fundo
asset_props = self.get_asset_properties(asset_key)
if asset_props.empty:
continue # Ativo nao eh valido
props = {
"asset_type" : asset_props["Type"], "group" : asset_props["Group"], "ticker_bbg" : asset_props["Ticker_BBG"],
"curncy_exp" : asset_props["Currency Exposure"], "location" : asset_props["Location"]
}
self.funds[fund_name].register_asset(asset_props["Asset"], asset_props["Ticker"], props, self.lq)
def process_trades(self, fund_names=None, date_reference=None):
if date_reference is None: date_reference = self.today
flag_modified = False
flag_reset = False
if fund_names is None: #Faz para todos os fundos cadastrados
fund_names = list(self.funds.keys())
# Nesse caso processa todos os fundos
for f in fund_names:
# Processa apenas para os fundos informados, se existirem
if f in self.funds:
# Primeiro criamos ativos novos (caso onde foram feitos trades com ativos novos)
self._register_new_assets(fund_name=f)
# Agora podemos processas os trades#
_flag_modified, _flag_reset = self.funds[f].process_trades(self.lq)
flag_modified = flag_modified | _flag_modified
flag_reset = flag_reset | _flag_reset
return flag_modified, flag_reset
def stop(self):
self.save_funds()
#self.export_data_to_db()
def update(self):
# Faz diferenca olhar se eh fim de semana ou nao? Acho que nao mais...
#if self.today != dt.date.today():
# tday= dt.date.today()
# if tday.weekday() <= 6:
# if not self.lq.is_holiday(tday, location="any"):
# self.source_bases_exec_count = 0
# self.today = tday
# else: return
# else: return #final de semana, nao vamos rodar
# Foi criado um script separado para atualizacao das source_bases
#if self.source_bases_exec_count == 0:
# now = dt.datetime.now()
# if now.hour*100 + now.minute >= 50:
# source_bases_updater.update(is_develop_mode=self.is_develop_mode)
# self.source_bases_exec_count += 1
# Processamos todos os trades
self.redirect_trades()
positions_changed, flag_reset = self.process_trades()
if not flag_reset:
if positions_changed:
# Enviamos as tabelas resultantes para a base de dados
self.export_data_to_db()
# Salvar pickle com informacao de processamento dos fundos
self.save_funds()
# Buscando alteracoes na historico_trades_ativos. Carregando caso haja.
self.dl.scan_files(export_to_blob=True)
# Atualizando tabelas carregadas.
self.lq.update()
else:
logger.warning("RESET automático agendado para próxima exec.")
self.__initialize(reset=flag_reset, is_develop_mode=self.is_develop_mode)
|