pentarosarium/processor
0
1import streamlit as st2import pandas as pd3import time4import matplotlib.pyplot as plt5from openpyxl.utils.dataframe import dataframe_to_rows6import io7from rapidfuzz import fuzz8import os9from openpyxl import load_workbook10from langchain.prompts import PromptTemplate11from langchain_core.runnables import RunnablePassthrough12from io import StringIO, BytesIO13import sys14import contextlib15from langchain_openai import ChatOpenAI # Updated import16import pdfkit17from jinja2 import Template18import time19from tenacity import retry, stop_after_attempt, wait_exponential20from typing import Optional21import torch22from transformers import (23 pipeline,24 AutoModelForSeq2SeqLM,25 AutoTokenizer,26 AutoModelForCausalLM # 4 Qwen27)28 29from threading import Event30import threading31from queue import Queue32 33from deep_translator import GoogleTranslator34from googletrans import Translator as LegacyTranslator35import plotly.graph_objects as go36from datetime import datetime37import plotly.express as px38 39 40class ProcessControl:41 def __init__(self):42 self.pause_event = Event()43 self.stop_event = Event()44 self.pause_event.set() # Start in non-paused state45 46 def pause(self):47 self.pause_event.clear()48 49 def resume(self):50 self.pause_event.set()51 52 def stop(self):53 self.stop_event.set()54 self.pause_event.set() # Ensure not stuck in pause55 56 def reset(self):57 self.stop_event.clear()58 self.pause_event.set()59 60 def is_paused(self):61 return not self.pause_event.is_set()62 63 def is_stopped(self):64 return self.stop_event.is_set()65 66 def wait_if_paused(self):67 self.pause_event.wait()68 69 70class FallbackLLMSystem:71 def __init__(self):72 """Initialize fallback models for event detection and reasoning"""73 try:74 # Initialize MT5 model (multilingual T5)75 self.model_name = "google/mt5-small"76 self.tokenizer = AutoTokenizer.from_pretrained(self.model_name)77 self.model = AutoModelForSeq2SeqLM.from_pretrained(self.model_name)78 79 # Set device80 self.device = "cuda" if torch.cuda.is_available() else "cpu"81 self.model = self.model.to(self.device)82 83 st.success(f"пока все в порядке: запущена MT5 model на = {self.device} =")84 85 except Exception as e:86 st.error(f"Ошибка запуска модели MT5: {str(e)}")87 raise88 89 def invoke(self, prompt_args):90 """Make the class compatible with LangChain by implementing invoke"""91 try:92 if isinstance(prompt_args, dict):93 # Extract the prompt template result94 template_result = prompt_args.get('template_result', '')95 if not template_result:96 # Try to construct from entity and news if available97 entity = prompt_args.get('entity', '')98 news = prompt_args.get('news', '')99 template_result = f"Analyze news about {entity}: {news}"100 else:101 template_result = str(prompt_args)102 103 # Process with MT5104 inputs = self.tokenizer(105 template_result,106 return_tensors="pt",107 padding=True,108 truncation=True,109 max_length=512110 ).to(self.device)111 112 outputs = self.model.generate(113 **inputs,114 max_length=200,115 num_return_sequences=1,116 do_sample=False,117 pad_token_id=self.tokenizer.pad_token_id118 )119 120 response = self.tokenizer.decode(outputs[0], skip_special_tokens=True)121 122 # Return in a format compatible with LangChain123 return type('Response', (), {'content': response})()124 125 except Exception as e:126 st.warning(f"MT5 generation error: {str(e)}")127 # Return a default response on error128 return type('Response', (), {129 'content': 'Impact: Неопределенный эффект\nReasoning: Ошибка анализа'130 })()131 132 def __or__(self, other):133 """Implement the | operator for chain compatibility"""134 if callable(other):135 return lambda x: other(self(x))136 return NotImplemented137 138 def __rrshift__(self, other):139 """Implement the >> operator for chain compatibility"""140 return self.__or__(other)141 142 def __call__(self, prompt_args):143 """Make the class callable for chain compatibility"""144 return self.invoke(prompt_args)145 146 def detect_events(self, text: str, entity: str) -> tuple[str, str]:147 """148 Detect events using MT5 with improved error handling and response parsing149 150 Args:151 text (str): The news text to analyze152 entity (str): The company/entity name153 154 Returns:155 tuple[str, str]: (event_type, summary)156 """157 # Initialize default return values158 event_type = "Нет"159 summary = ""160 161 # Input validation162 if not text or not entity or not isinstance(text, str) or not isinstance(entity, str):163 return event_type, "Invalid input"164 165 try:166 # Clean and prepare input text167 text = text.strip()168 entity = entity.strip()169 170 # Construct prompt with better formatting171 prompt = f"""<s>Analyze the following news about {entity}:172 173 Text: {text}174 175 Task: Identify the main event type and provide a brief summary.176 177 Event types:178 1. Отчетность - Events related to financial reports, earnings, revenue, EBITDA179 2. РЦБ - Events related to securities, bonds, stock market, defaults, restructuring180 3. Суд - Events related to legal proceedings, lawsuits, arbitration181 4. Нет - No significant events detected182 183 Required output format:184 Тип: [event type]185 Краткое описание: [1-2 sentence summary]</s>"""186 187 # Process with MT5188 try:189 inputs = self.tokenizer(190 prompt,191 return_tensors="pt",192 padding=True,193 truncation=True,194 max_length=512195 ).to(self.device)196 197 outputs = self.model.generate(198 **inputs,199 max_length=300, # Increased for better summaries200 num_return_sequences=1,201 do_sample=False,202 pad_token_id=self.tokenizer.pad_token_id,203 eos_token_id=self.tokenizer.eos_token_id,204 no_repeat_ngram_size=3 # Prevent repetition205 )206 207 response = self.tokenizer.decode(outputs[0], skip_special_tokens=True)208 209 except torch.cuda.OutOfMemoryError:210 st.warning("GPU memory exceeded, falling back to CPU")211 self.model = self.model.to('cpu')212 inputs = inputs.to('cpu')213 outputs = self.model.generate(214 **inputs,215 max_length=300,216 num_return_sequences=1,217 do_sample=False,218 pad_token_id=self.tokenizer.pad_token_id219 )220 response = self.tokenizer.decode(outputs[0], skip_special_tokens=True)221 self.model = self.model.to(self.device) # Move back to GPU222 223 # Enhanced response parsing224 if "Тип:" in response and "Краткое описание:" in response:225 try:226 # Split and clean parts227 parts = response.split("Краткое описание:")228 type_part = parts[0].split("Тип:")[1].strip()229 230 # Validate event type with fuzzy matching231 valid_types = ["Отчетность", "РЦБ", "Суд", "Нет"]232 233 # Check for exact matches first234 if type_part in valid_types:235 event_type = type_part236 else:237 # Check keywords for each type238 keywords = {239 "Отчетность": ["отчет", "выручка", "прибыль", "ebitda", "финанс"],240 "РЦБ": ["облигаци", "купон", "дефолт", "реструктуризац", "ценные бумаги"],241 "Суд": ["суд", "иск", "арбитраж", "разбирательств"]242 }243 244 # Look for keywords in both type and summary245 full_text = response.lower()246 for event_category, category_keywords in keywords.items():247 if any(keyword in full_text for keyword in category_keywords):248 event_type = event_category249 break250 251 # Extract and clean summary252 if len(parts) > 1:253 summary = parts[1].strip()254 # Ensure summary isn't too long255 if len(summary) > 200:256 summary = summary[:197] + "..."257 258 # Add entity reference if missing259 if entity.lower() not in summary.lower():260 summary = f"Компания {entity}: {summary}"261 262 except IndexError:263 st.warning("Error parsing model response format")264 return "Нет", "Error parsing response"265 266 # Additional validation267 if not summary or len(summary) < 5:268 keywords = {269 "Отчетность": "Обнаружена информация о финансовой отчетности",270 "РЦБ": "Обнаружена информация о ценных бумагах",271 "Суд": "Обнаружена информация о судебном разбирательстве",272 "Нет": "Значимых событий не обнаружено"273 }274 summary = f"{keywords.get(event_type, 'Требуется дополнительный анализ')} ({entity})"275 276 return event_type, summary277 278 except Exception as e:279 st.warning(f"Event detection error: {str(e)}")280 # Try to provide more specific error information281 if "CUDA" in str(e):282 return "Нет", "GPU error - falling back to CPU needed"283 elif "tokenizer" in str(e):284 return "Нет", "Text processing error"285 elif "model" in str(e):286 return "Нет", "Model inference error"287 else:288 return "Нет", "Ошибка анализа"289 290 291def ensure_groq_llm():292 """Initialize Groq LLM for impact estimation"""293 try:294 if 'groq_key' not in st.secrets:295 st.error("Groq API key not found in secrets. Please add it with the key 'groq_key'.")296 return None297 298 return ChatOpenAI(299 base_url="https://api.groq.com/openai/v1",300 model="llama-3.1-70b-versatile",301 openai_api_key=st.secrets['groq_key'],302 temperature=0.0303 )304 except Exception as e:305 st.error(f"Error initializing Groq LLM: {str(e)}")306 return None307 308def estimate_impact(llm, news_text, entity):309 """310 Estimate impact using Groq LLM regardless of the main model choice.311 Falls back to the provided LLM if Groq initialization fails.312 """313 # Initialize default return values314 impact = "Неопределенный эффект"315 reasoning = "Не удалось получить обоснование"316 317 try:318 # Always try to use Groq first319 groq_llm = ensure_groq_llm()320 working_llm = groq_llm if groq_llm is not None else llm321 322 template = """323 You are a financial analyst. Analyze this news piece about {entity} and assess its potential impact.324 325 News: {news}326 327 Classify the impact into one of these categories:328 1. "Значительный риск убытков" (Significant loss risk)329 2. "Умеренный риск убытков" (Moderate loss risk)330 3. "Незначительный риск убытков" (Minor loss risk)331 4. "Вероятность прибыли" (Potential profit)332 5. "Неопределенный эффект" (Uncertain effect)333 334 Provide a brief, fact-based reasoning for your assessment.335 336 Format your response exactly as:337 Impact: [category]338 Reasoning: [explanation in 2-3 sentences]339 """340 341 prompt = PromptTemplate(template=template, input_variables=["entity", "news"])342 chain = prompt | working_llm343 response = chain.invoke({"entity": entity, "news": news_text})344 345 # Extract content from response346 response_text = response.content if hasattr(response, 'content') else str(response)347 348 if "Impact:" in response_text and "Reasoning:" in response_text:349 impact_part, reasoning_part = response_text.split("Reasoning:")350 impact_temp = impact_part.split("Impact:")[1].strip()351 352 # Validate impact category353 valid_impacts = [354 "Значительный риск убытков",355 "Умеренный риск убытков",356 "Незначительный риск убытков",357 "Вероятность прибыли",358 "Неопределенный эффект"359 ]360 if impact_temp in valid_impacts:361 impact = impact_temp362 reasoning = reasoning_part.strip()363 364 except Exception as e:365 st.warning(f"Error in impact estimation: {str(e)}")366 367 return impact, reasoning368 369class QwenSystem:370 def __init__(self):371 """Initialize Qwen 2.5 Coder model"""372 try:373 self.model_name = "Qwen/Qwen2.5-Coder-32B-Instruct"374 375 # Initialize model with auto settings376 self.model = AutoModelForCausalLM.from_pretrained(377 self.model_name,378 torch_dtype="auto",379 device_map="auto"380 )381 self.tokenizer = AutoTokenizer.from_pretrained(self.model_name)382 383 st.success(f"запустил Qwen2.5 model")384 385 except Exception as e:386 st.error(f"ошибка запуска Qwen2.5: {str(e)}")387 raise388 389 def invoke(self, messages):390 """Process messages using Qwen's chat template"""391 try:392 # Prepare messages with system prompt393 chat_messages = [394 {"role": "system", "content": "You are wise financial analyst. You are a helpful assistant."}395 ]396 chat_messages.extend(messages)397 398 # Apply chat template399 text = self.tokenizer.apply_chat_template(400 chat_messages,401 tokenize=False,402 add_generation_prompt=True403 )404 405 # Prepare model inputs406 model_inputs = self.tokenizer([text], return_tensors="pt").to(self.model.device)407 408 # Generate response409 generated_ids = self.model.generate(410 **model_inputs,411 max_new_tokens=512,412 pad_token_id=self.tokenizer.pad_token_id,413 eos_token_id=self.tokenizer.eos_token_id414 )415 416 # Extract new tokens417 generated_ids = [418 output_ids[len(input_ids):] 419 for input_ids, output_ids in zip(model_inputs.input_ids, generated_ids)420 ]421 422 # Decode response423 response = self.tokenizer.batch_decode(generated_ids, skip_special_tokens=True)[0]424 425 # Return in ChatOpenAI-compatible format426 return type('Response', (), {'content': response})()427 428 except Exception as e:429 st.warning(f"Qwen generation error: {str(e)}")430 raise431 432 433class ProcessingUI:434 def __init__(self):435 if 'control' not in st.session_state:436 st.session_state.control = ProcessControl()437 438 # Initialize processing stats in session state if not exists439 if 'processing_stats' not in st.session_state:440 st.session_state.processing_stats = {441 'start_time': time.time(),442 'entities': {},443 'events_timeline': [],444 'negative_alerts': [],445 'processing_speed': []446 }447 448 # Create main layout449 self.setup_layout()450 451 def setup_layout(self):452 """Setup the main UI layout with tabs and sections"""453 # Control Panel454 with st.container():455 col1, col2, col3 = st.columns([2,2,1])456 with col1:457 if st.button(458 "⏸️ Пауза" if not st.session_state.control.is_paused() else "▶️ Продолжить",459 use_container_width=True460 ):461 if st.session_state.control.is_paused():462 st.session_state.control.resume()463 else:464 st.session_state.control.pause()465 with col2:466 if st.button("⏹️ Остановить", use_container_width=True):467 st.session_state.control.stop()468 with col3:469 self.timer_display = st.empty()470 471 # Progress Bar with custom styling472 st.markdown("""473 <style>474 .stProgress > div > div > div > div {475 background-image: linear-gradient(to right, #FF6B6B, #4ECDC4);476 }477 </style>""", 478 unsafe_allow_html=True479 )480 self.progress_bar = st.progress(0)481 self.status = st.empty()482 483 # Create tabs for different views484 tab1, tab2, tab3, tab4 = st.tabs([485 "📊 Основные метрики", 486 "🏢 По организациям", 487 "⚠️ Важные события", 488 "📈 Аналитика"489 ])490 491 with tab1:492 self.setup_main_metrics_tab()493 494 with tab2:495 self.setup_entity_tab()496 497 with tab3:498 self.setup_events_tab()499 500 with tab4:501 self.setup_analytics_tab()502 503 def setup_entity_tab(self):504 """Setup the entity-wise analysis display"""505 # Entity filter506 self.entity_filter = st.multiselect(507 "Фильтр по организациям:",508 options=[], # Will be populated as entities are processed509 default=None510 )511 512 # Entity metrics513 self.entity_cols = st.columns([2,1,1,1])514 self.entity_chart = st.empty()515 self.entity_table = st.empty()516 517 def setup_events_tab(self):518 """Setup the events timeline display"""519 # Event type filter - store in session state520 if 'event_filter' not in st.session_state:521 st.session_state.event_filter = []522 523 st.session_state.event_filter = st.multiselect(524 "Тип события:",525 options=["Отчетность", "РЦБ", "Суд"],526 default=None,527 key="event_filter_key"528 )529 530 self.timeline_container = st.container()531 532 def _update_events_view(self, row, event_type):533 """Update events timeline"""534 if event_type != 'Нет':535 event_html = f"""536 <div class='timeline-item' style='537 border-left: 4px solid #2196F3;538 margin: 10px 0;539 padding: 10px;540 background: #f5f5f5;541 border-radius: 4px;542 '>543 <h4 style='color: #2196F3; margin: 0;'>{event_type}</h4>544 <p><strong>{row['Объект']}</strong></p>545 <p>{row['Заголовок']}</p>546 <p style='font-size: 0.9em;'>{row['Выдержки из текста']}</p>547 <small style='color: #666;'>{datetime.now().strftime('%H:%M:%S')}</small>548 </div>549 """550 with self.timeline_container:551 st.markdown(event_html, unsafe_allow_html=True)552 553 def setup_analytics_tab(self):554 """Setup the analytics display"""555 # Create containers for analytics556 self.speed_container = st.container()557 with self.speed_container:558 st.subheader("Скорость обработки")559 self.speed_chart = st.empty()560 561 self.sentiment_container = st.container()562 with self.sentiment_container:563 st.subheader("Распределение тональности")564 self.sentiment_chart = st.empty()565 566 self.correlation_container = st.container()567 with self.correlation_container:568 st.subheader("Корреляция между метриками")569 self.correlation_chart = st.empty()570 571 def update_stats(self, row, sentiment, event_type, processing_speed):572 """Update all statistics and displays"""573 # Update session state stats574 stats = st.session_state.processing_stats575 entity = row['Объект']576 577 # Update entity stats578 if entity not in stats['entities']:579 stats['entities'][entity] = {580 'total': 0,581 'negative': 0,582 'events': 0,583 'timeline': []584 }585 586 stats['entities'][entity]['total'] += 1587 if sentiment == 'Negative':588 stats['entities'][entity]['negative'] += 1589 if event_type != 'Нет':590 stats['entities'][entity]['events'] += 1591 592 # Update processing speed593 stats['processing_speed'].append(processing_speed)594 595 # Update UI components596 self._update_main_metrics(row, sentiment, event_type, processing_speed)597 self._update_entity_view()598 self._update_events_view(row, event_type)599 self._update_analytics()600 601 def _update_main_metrics(self, row, sentiment, event_type, speed):602 """Update main metrics tab"""603 total = sum(e['total'] for e in st.session_state.processing_stats['entities'].values())604 total_negative = sum(e['negative'] for e in st.session_state.processing_stats['entities'].values())605 total_events = sum(e['events'] for e in st.session_state.processing_stats['entities'].values())606 607 # Update metrics608 self.total_processed.metric("Обработано", total)609 self.negative_count.metric("Негативных", total_negative)610 self.events_count.metric("Событий", total_events)611 self.speed_metric.metric("Скорость", f"{speed:.1f} сообщ/сек")612 613 # Update recent items614 self._update_recent_items(row, sentiment, event_type)615 616 def _update_recent_items(self, row, sentiment, event_type):617 """Update recent items display using Streamlit native components"""618 if 'recent_items' not in st.session_state:619 st.session_state.recent_items = []620 621 # Add new item to the list622 new_item = {623 'entity': row['Объект'],624 'headline': row['Заголовок'],625 'sentiment': sentiment,626 'event_type': event_type,627 'time': datetime.now().strftime('%H:%M:%S')628 }629 630 # Update the list in session state631 if not any(632 item['entity'] == new_item['entity'] and 633 item['headline'] == new_item['headline'] 634 for item in st.session_state.recent_items635 ):636 st.session_state.recent_items.insert(0, new_item)637 st.session_state.recent_items = st.session_state.recent_items[:10] # Keep last 10 items638 639 # Prepare markdown for all items640 all_items_markdown = ""641 642 for item in st.session_state.recent_items:643 if item['sentiment'] in ['Positive', 'Negative']:644 sentiment_color = "🔴" if item['sentiment'] == 'Negative' else "🟢"645 event_icon = "📅" if item['event_type'] != 'Нет' else ""646 647 event_text = f" | Событие: {item['event_type']}" if item['event_type'] != 'Нет' else ""648 649 all_items_markdown += f"""650 {sentiment_color} **{item['entity']}** {event_icon}651 652 {item['headline']}653 654 *{item['sentiment']}*{event_text} | {item['time']}655 656 ---657 """658 659 # Update container with all items at once660 if all_items_markdown:661 self.recent_items_container.markdown(all_items_markdown)662 663 def setup_main_metrics_tab(self):664 """Setup the main metrics display with updated styling"""665 # Create metrics containers666 metrics_cols = st.columns(4)667 self.total_processed = metrics_cols[0].empty()668 self.negative_count = metrics_cols[1].empty()669 self.events_count = metrics_cols[2].empty()670 self.speed_metric = metrics_cols[3].empty()671 672 # Create container for recent items673 st.markdown("### негативные/позитивные")674 self.recent_items_container = st.empty()675 676 677 def _update_entity_view(self):678 """Update entity tab visualizations"""679 stats = st.session_state.processing_stats['entities']680 if not stats:681 return682 683 # Get filtered entities684 filtered_entities = self.entity_filter or stats.keys()685 686 # Create entity comparison chart using Plotly687 df_entities = pd.DataFrame.from_dict(stats, orient='index')688 df_entities = df_entities.loc[filtered_entities] # Apply filter689 690 fig = go.Figure(data=[691 go.Bar(692 name='Всего',693 x=df_entities.index,694 y=df_entities['total'],695 marker_color='#E0E0E0' # Light gray696 ),697 go.Bar(698 name='Негативные',699 x=df_entities.index,700 y=df_entities['negative'],701 marker_color='#FF6B6B' # Red702 ),703 go.Bar(704 name='События',705 x=df_entities.index,706 y=df_entities['events'],707 marker_color='#2196F3' # Blue708 )709 ])710 711 fig.update_layout(712 barmode='group',713 title='Статистика по организациям',714 xaxis_title='Организация',715 yaxis_title='Количество',716 showlegend=True717 )718 719 self.entity_chart.plotly_chart(fig, use_container_width=True)720 721 def _update_analytics(self):722 """Update analytics tab visualizations"""723 stats = st.session_state.processing_stats724 725 # Processing speed chart - showing last 20 measurements726 speeds = stats['processing_speed'][-20:]727 if speeds:728 fig_speed = go.Figure(data=go.Scatter(729 y=speeds,730 mode='lines+markers',731 name='Скорость',732 line=dict(color='#4CAF50')733 ))734 fig_speed.update_layout(735 title='Скорость обработки',736 yaxis_title='Сообщений в секунду',737 showlegend=True738 )739 self.speed_chart.plotly_chart(fig_speed, use_container_width=True)740 741 # Sentiment distribution pie chart742 if stats['entities']:743 total_negative = sum(e['negative'] for e in stats['entities'].values())744 total_positive = sum(e['events'] for e in stats['entities'].values())745 total_neutral = sum(e['total'] for e in stats['entities'].values()) - total_negative - total_positive746 747 fig_sentiment = go.Figure(data=[go.Pie(748 labels=['Негативные', 'Позитивные', 'Нейтральные'],749 values=[total_negative, total_positive, total_neutral],750 marker_colors=['#FF6B6B', '#4ECDC4', '#95A5A6']751 )])752 self.sentiment_chart.plotly_chart(fig_sentiment, use_container_width=True)753 754 def update_progress(self, current, total):755 """Update progress bar, elapsed time and estimated time remaining"""756 progress = current / total757 self.progress_bar.progress(progress)758 self.status.text(f"Обрабатываем {current} из {total} сообщений...")759 760 # Calculate times761 current_time = time.time()762 elapsed = current_time - st.session_state.processing_stats['start_time']763 764 # Calculate processing speed and estimated time remaining765 if current > 0:766 speed = current / elapsed # items per second767 remaining_items = total - current768 estimated_remaining = remaining_items / speed if speed > 0 else 0769 770 time_display = (771 f"⏱️ Прошло: {format_elapsed_time(elapsed)} | "772 f"Осталось: {format_elapsed_time(estimated_remaining)}"773 )774 else:775 time_display = f"⏱️ Прошло: {format_elapsed_time(elapsed)}"776 777 self.timer_display.markdown(time_display)778 779 780class EventDetectionSystem:781 def __init__(self):782 try:783 # Initialize models with specific labels784 self.finbert = pipeline(785 "text-classification", 786 model="ProsusAI/finbert",787 return_all_scores=True788 )789 self.business_classifier = pipeline(790 "text-classification", 791 model="yiyanghkust/finbert-tone",792 return_all_scores=True793 )794 st.success("продолжается пока хорошо: BERT-модели запущены для детекции новостей")795 except Exception as e:796 st.error(f"Ошибка запуска BERT: {str(e)}")797 raise798 799 def detect_event_type(self, text, entity):800 event_type = "Нет"801 summary = ""802 803 try:804 # Ensure text is properly formatted805 text = str(text).strip()806 if not text:807 return "Нет", "Empty text"808 809 # Get predictions810 finbert_scores = self.finbert(811 text,812 truncation=True,813 max_length=512814 )815 business_scores = self.business_classifier(816 text,817 truncation=True,818 max_length=512819 )820 821 # Get highest scoring predictions822 finbert_pred = max(finbert_scores[0], key=lambda x: x['score'])823 business_pred = max(business_scores[0], key=lambda x: x['score'])824 825 # Map to event types with confidence threshold826 confidence_threshold = 0.6827 max_confidence = max(finbert_pred['score'], business_pred['score'])828 829 if max_confidence >= confidence_threshold:830 if any(term in text.lower() for term in ['отчет', 'выручка', 'прибыль', 'ebitda']):831 event_type = "Отчетность"832 summary = f"Финансовая отчетность (confidence: {max_confidence:.2f})"833 elif any(term in text.lower() for term in ['облигаци', 'купон', 'дефолт', 'реструктуризац']):834 event_type = "РЦБ"835 summary = f"Событие РЦБ (confidence: {max_confidence:.2f})"836 elif any(term in text.lower() for term in ['суд', 'иск', 'арбитраж']):837 event_type = "Суд"838 summary = f"Судебное разбирательство (confidence: {max_confidence:.2f})"839 840 if event_type != "Нет":841 summary += f"\nКомпания: {entity}"842 843 return event_type, summary844 845 except Exception as e:846 st.warning(f"Event detection error: {str(e)}")847 return "Нет", "Error in event detection"848 849class TranslationSystem:850 def __init__(self):851 """Initialize translation system using Helsinki NLP model with fallback options"""852 try:853 self.translator = pipeline("translation", model="Helsinki-NLP/opus-mt-ru-en")854 # Initialize fallback translator855 self.fallback_translator = GoogleTranslator(source='ru', target='en')856 self.legacy_translator = LegacyTranslator()857 st.success("начинается все хорошо: запустил систему перевода")858 except Exception as e:859 st.error(f"Ошибка запуска перевода: {str(e)}")860 raise861 862 def _split_into_chunks(self, text: str, max_length: int = 450) -> list:863 """Split text into chunks while preserving word boundaries"""864 words = text.split()865 chunks = []866 current_chunk = []867 current_length = 0868 869 for word in words:870 word_length = len(word)871 if current_length + word_length + 1 <= max_length:872 current_chunk.append(word)873 current_length += word_length + 1874 else:875 if current_chunk:876 chunks.append(' '.join(current_chunk))877 current_chunk = [word]878 current_length = word_length879 880 if current_chunk:881 chunks.append(' '.join(current_chunk))882 883 return chunks884 885 def _translate_chunk_with_retries(self, chunk: str, max_retries: int = 3) -> str:886 """Attempt translation with multiple fallback options"""887 if not chunk or not chunk.strip():888 return ""889 890 for attempt in range(max_retries):891 try:892 # First try Helsinki NLP893 result = self.translator(chunk, max_length=512)894 if result and isinstance(result, list) and len(result) > 0:895 translated = result[0].get('translation_text')896 if translated and isinstance(translated, str):897 return translated898 899 # First fallback: Google Translator900 translated = self.fallback_translator.translate(chunk)901 if translated and isinstance(translated, str):902 return translated903 904 # Second fallback: Legacy Google Translator905 translated = self.legacy_translator.translate(chunk, src='ru', dest='en').text906 if translated and isinstance(translated, str):907 return translated908 909 except Exception as e:910 if attempt == max_retries - 1:911 st.warning(f"Попробовал перевести {max_retries} раз, не преуспел: {str(e)}")912 time.sleep(1 * (attempt + 1)) # Exponential backoff913 914 return chunk # Return original text if all translation attempts fail915 916 def translate_text(self, text: str) -> str:917 """Translate text with robust error handling and validation"""918 # Input validation919 if pd.isna(text) or not isinstance(text, str):920 return str(text) if pd.notna(text) else ""921 922 text = str(text).strip()923 if not text:924 return ""925 926 try:927 # Split into manageable chunks928 chunks = self._split_into_chunks(text)929 translated_chunks = []930 931 # Process each chunk with validation932 for chunk in chunks:933 if not chunk.strip():934 continue935 936 translated_chunk = self._translate_chunk_with_retries(chunk)937 if translated_chunk: # Only add non-empty translations938 translated_chunks.append(translated_chunk)939 time.sleep(0.1) # Rate limiting940 941 # Final validation of results942 if not translated_chunks:943 return text # Return original if no translations succeeded944 945 result = ' '.join(translated_chunks)946 return result if result.strip() else text947 948 except Exception as e:949 st.warning(f"Translation error: {str(e)}")950 return text # Return original text on error951 952 953 954def process_file(uploaded_file, model_choice, translation_method=None):955 df = None956 processed_rows_df = pd.DataFrame()957 last_time = time.time()958 959 try:960 # Initialize UI and control systems961 ui = ProcessingUI()962 translator = TranslationSystem()963 event_detector = EventDetectionSystem()964 965 # Load and prepare data966 df = pd.read_excel(uploaded_file, sheet_name='Публикации')967 llm = init_langchain_llm(model_choice)968 969 # Initialize Groq for impact estimation970 groq_llm = ensure_groq_llm()971 if groq_llm is None:972 st.warning("Failed to initialize Groq LLM for impact estimation. Using fallback model.")973 974 # Initialize all required columns at the start975 required_columns = {976 'Объект': '',977 'Заголовок': '',978 'Выдержки из текста': '',979 'Translated': '',980 'Sentiment': 'Neutral',981 'Impact': 'Неопределенный эффект',982 'Reasoning': 'Не проанализировано',983 'Event_Type': 'Нет',984 'Event_Summary': ''985 }986 987 # Ensure all required columns exist in DataFrame988 for col, default_value in required_columns.items():989 if col not in df.columns:990 df[col] = default_value991 992 # Create processed_rows_df with all columns from original df and required columns993 all_columns = list(set(list(df.columns) + list(required_columns.keys())))994 processed_rows_df = pd.DataFrame(columns=all_columns)995 996 # Deduplication997 original_count = len(df)998 df = df.groupby('Объект', group_keys=False).apply(999 lambda x: fuzzy_deduplicate(x, 'Выдержки из текста', 55)1000 ).reset_index(drop=True)1001 st.write(f"Из {original_count} сообщений удалено {original_count - len(df)} дубликатов.")1002 1003 # Process rows1004 total_rows = len(df)1005 processed_rows = 01006 grlm = init_langchain_llm("Groq (llama-3.1-70b)")1007 1008 for idx, row in df.iterrows():1009 if st.session_state.control.is_stopped():1010 st.warning("Обработку остановили")1011 if not processed_rows_df.empty:1012 try:1013 # Create the output files for each sheet1014 monitoring_df = processed_rows_df[processed_rows_df['Event_Type'] != 'Нет'].copy()1015 svodka_df = processed_rows_df.groupby('Объект').agg({1016 'Объект': 'first',1017 'Sentiment': lambda x: sum(x == 'Negative'),1018 'Event_Type': lambda x: sum(x != 'Нет')1019 }).reset_index()1020 1021 # Prepare final DataFrame for file creation1022 result_df = pd.DataFrame()1023 result_df['Мониторинг'] = monitoring_df.to_dict('records')1024 result_df['Сводка'] = svodka_df.to_dict('records')1025 result_df['Публикации'] = processed_rows_df.to_dict('records')1026 1027 output = create_output_file(result_df, uploaded_file)1028 if output is not None:1029 st.download_button(1030 label=f"📊 Скачать результат ({processed_rows} из {total_rows} строк)",1031 data=output,1032 file_name="partial_analysis.xlsx",1033 mime="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",1034 key="partial_download"1035 )1036 except Exception as e:1037 st.error(f"Ошибка при создании файла: {str(e)}")1038 1039 return processed_rows_df1040 1041 st.session_state.control.wait_if_paused()1042 if st.session_state.control.is_paused():1043 continue1044 1045 try:1046 # Copy original row data1047 new_row = row.copy()1048 1049 # Translation1050 translated_text = translator.translate_text(row['Выдержки из текста'])1051 new_row['Translated'] = translated_text1052 1053 # Sentiment analysis1054 sentiment = analyze_sentiment(translated_text)1055 new_row['Sentiment'] = sentiment1056 1057 # Event detection1058 event_type, event_summary = event_detector.detect_event_type(1059 row['Выдержки из текста'],1060 row['Объект']1061 )1062 new_row['Event_Type'] = event_type1063 new_row['Event_Summary'] = event_summary1064 1065 # Handle negative sentiment1066 if sentiment == "Negative":1067 try:1068 if translated_text and len(translated_text.strip()) > 0:1069 impact, reasoning = estimate_impact(1070 groq_llm if groq_llm is not None else llm,1071 translated_text,1072 row['Объект']1073 )1074 new_row['Impact'] = impact1075 new_row['Reasoning'] = translate_reasoning_to_russian(grlm, reasoning)1076 except Exception as e:1077 new_row['Impact'] = "Неопределенный эффект"1078 new_row['Reasoning'] = "Ошибка анализа"1079 1080 # Add processed row to DataFrame1081 processed_rows_df = pd.concat([processed_rows_df, pd.DataFrame([new_row])], ignore_index=True)1082 1083 # Calculate processing speed1084 current_time = time.time()1085 processing_speed = 1.0 / (current_time - last_time) if (current_time - last_time) > 0 else 01086 last_time = current_time1087 1088 # Update UI stats1089 ui.update_stats(1090 row=new_row,1091 sentiment=sentiment,1092 event_type=event_type,1093 processing_speed=processing_speed1094 )1095 1096 # Update progress1097 processed_rows += 11098 ui.update_progress(processed_rows, total_rows)1099 1100 except Exception as e:1101 st.warning(f"Ошибка в обработке ряда {idx + 1}: {str(e)}")1102 continue1103 1104 return processed_rows_df1105 1106 except Exception as e:1107 st.error(f"Ошибка в обработке файла: {str(e)}")1108 return None1109 1110 1111 1112 1113def create_download_section(excel_data, pdf_data):1114 st.markdown("""1115 <div class="download-container">1116 <div class="download-header">📥 Результаты анализа доступны для скачивания:</div>1117 </div>1118 """, unsafe_allow_html=True)1119 1120 col1, col2 = st.columns(2)1121 1122 with col1:1123 if excel_data is not None:1124 st.download_button(1125 label="📊 Скачать Excel отчет",1126 data=excel_data,1127 file_name="результат_анализа.xlsx",1128 mime="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",1129 key="excel_download"1130 )1131 else:1132 st.error("Ошибка при создании Excel файла")1133 1134 1135 1136 1137def display_sentiment_results(row, sentiment, impact=None, reasoning=None):1138 if sentiment == "Negative":1139 st.markdown(f"""1140 <div style='color: red; font-weight: bold;'>1141 Объект: {row['Объект']}<br>1142 Новость: {row['Заголовок']}<br>1143 Тональность: {sentiment}<br>1144 {"Эффект: " + impact + "<br>" if impact else ""}1145 {"Обоснование: " + reasoning + "<br>" if reasoning else ""}1146 </div>1147 """, unsafe_allow_html=True)1148 elif sentiment == "Positive":1149 st.markdown(f"""1150 <div style='color: green; font-weight: bold;'>1151 Объект: {row['Объект']}<br>1152 Новость: {row['Заголовок']}<br>1153 Тональность: {sentiment}<br>1154 </div>1155 """, unsafe_allow_html=True)1156 else:1157 st.write(f"Объект: {row['Объект']}")1158 st.write(f"Новость: {row['Заголовок']}")1159 st.write(f"Тональность: {sentiment}")1160 1161 st.write("---")1162 1163 1164 1165 1166 1167# Initialize sentiment analyzers1168finbert = pipeline("sentiment-analysis", model="ProsusAI/finbert")1169roberta = pipeline("sentiment-analysis", model="cardiffnlp/twitter-roberta-base-sentiment")1170finbert_tone = pipeline("sentiment-analysis", model="yiyanghkust/finbert-tone")1171 1172 1173def get_mapped_sentiment(result):1174 label = result['label'].lower()1175 if label in ["positive", "label_2", "pos", "pos_label"]:1176 return "Positive"1177 elif label in ["negative", "label_0", "neg", "neg_label"]:1178 return "Negative"1179 return "Neutral"1180 1181 1182 1183def analyze_sentiment(text):1184 try:1185 finbert_result = get_mapped_sentiment(1186 finbert(text, truncation=True, max_length=512)[0]1187 )1188 roberta_result = get_mapped_sentiment(1189 roberta(text, truncation=True, max_length=512)[0]1190 )1191 finbert_tone_result = get_mapped_sentiment(1192 finbert_tone(text, truncation=True, max_length=512)[0]1193 )1194 1195 # Count occurrences of each sentiment1196 sentiments = [finbert_result, roberta_result, finbert_tone_result]1197 sentiment_counts = {s: sentiments.count(s) for s in set(sentiments)}1198 1199 # Return sentiment if at least two models agree1200 for sentiment, count in sentiment_counts.items():