diff --git a/apps/epg/tasks.py b/apps/epg/tasks.py index 9153485f..f605abc6 100644 --- a/apps/epg/tasks.py +++ b/apps/epg/tasks.py @@ -532,11 +532,20 @@ def parse_channels_only(source): source.last_message = f"No URL provided, cannot fetch EPG data" source.save(update_fields=['updated_at']) - # Initialize process variable for memory tracking + # Initialize process variable for memory tracking only in debug mode try: - process = psutil.Process() - initial_memory = process.memory_info().rss / 1024 / 1024 - logger.info(f"[parse_channels_only] Initial memory usage: {initial_memory:.2f} MB") + process = None + # Get current log level as a number rather than string + current_log_level = logger.getEffectiveLevel() + + # Only track memory usage when log level is DEBUG (10) or more verbose + # This is more future-proof than string comparisons + if current_log_level <= logging.DEBUG or settings.DEBUG: + process = psutil.Process() + initial_memory = process.memory_info().rss / 1024 / 1024 + logger.debug(f"[parse_channels_only] Initial memory usage: {initial_memory:.2f} MB") + else: + logger.debug("Memory tracking disabled in production mode") except (ImportError, NameError): process = None logger.warning("psutil not available for memory tracking") @@ -596,230 +605,219 @@ def parse_channels_only(source): total_channels = 500 # Default estimate # Close the file to reset position - logger.info(f"Closing initial file handle") + logger.debug(f"Closing initial file handle") source_file.close() if process: - logger.info(f"[parse_channels_only] Memory after closing initial file: {process.memory_info().rss / 1024 / 1024:.2f} MB") + logger.debug(f"[parse_channels_only] Memory after closing initial file: {process.memory_info().rss / 1024 / 1024:.2f} MB") # Update progress after counting send_epg_update(source.id, "parsing_channels", 25, total_channels=total_channels) # Reset file position for actual processing - logger.info(f"Re-opening file for channel parsing: {file_path}") + logger.debug(f"Re-opening file for channel parsing: {file_path}") source_file = gzip.open(file_path, 'rb') if is_gzipped else open(file_path, 'rb') if process: - logger.info(f"[parse_channels_only] Memory after re-opening file: {process.memory_info().rss / 1024 / 1024:.2f} MB") + logger.debug(f"[parse_channels_only] Memory after re-opening file: {process.memory_info().rss / 1024 / 1024:.2f} MB") - logger.info(f"Creating iterparse context") - channel_parser = etree.iterparse(source_file, events=('end',), tag='channel') + # Change iterparse to look for both channel and programme elements + logger.debug(f"Creating iterparse context for channels and programmes") + channel_parser = etree.iterparse(source_file, events=('end',), tag=('channel', 'programme')) if process: - logger.info(f"[parse_channels_only] Memory after creating iterparse: {process.memory_info().rss / 1024 / 1024:.2f} MB") + logger.debug(f"[parse_channels_only] Memory after creating iterparse: {process.memory_info().rss / 1024 / 1024:.2f} MB") channel_count = 0 total_elements_processed = 0 # Track total elements processed, not just channels for _, elem in channel_parser: total_elements_processed += 1 - channel_count += 1 - tvg_id = elem.get('id', '').strip() - if tvg_id: - display_name = None - for child in elem: - if child.tag == 'display-name' and child.text: - display_name = child.text.strip() - break - if not display_name: - display_name = tvg_id - - # Use lazy loading approach to reduce memory usage - if tvg_id in existing_tvg_ids: - # Only fetch the object if we need to update it and it hasn't been loaded yet - if tvg_id not in existing_epgs: - try: - # This loads the full EPG object from the database and caches it - existing_epgs[tvg_id] = EPGData.objects.get(tvg_id=tvg_id, epg_source=source) - except EPGData.DoesNotExist: - # Handle race condition where record was deleted - existing_tvg_ids.remove(tvg_id) - epgs_to_create.append(EPGData( - tvg_id=tvg_id, - name=display_name, - epg_source=source, - )) - logger.debug(f"[parse_channels_only] Added new channel to epgs_to_create 1: {tvg_id} - {display_name}") - continue - - # We use the cached object to check if the name has changed - epg_obj = existing_epgs[tvg_id] - if epg_obj.name != display_name: - # Only update if the name actually changed - epg_obj.name = display_name - epgs_to_update.append(epg_obj) - logger.debug(f"[parse_channels_only] Added channel to update to epgs_to_update: {tvg_id} - {display_name}") - else: - # This is a new channel that doesn't exist in our database - epgs_to_create.append(EPGData( - tvg_id=tvg_id, - name=display_name, - epg_source=source, - )) - logger.debug(f"[parse_channels_only] Added new channel to epgs_to_create 2: {tvg_id} - {display_name}") - - processed_channels += 1 - - # Batch processing - if len(epgs_to_create) >= batch_size: - logger.info(f"[parse_channels_only] Bulk creating {len(epgs_to_create)} EPG entries") - EPGData.objects.bulk_create(epgs_to_create, ignore_conflicts=True) - if process: - logger.info(f"[parse_channels_only] Memory after bulk_create: {process.memory_info().rss / 1024 / 1024:.2f} MB") - del epgs_to_create # Explicit deletion - epgs_to_create = [] - gc.collect() - if process: - logger.info(f"[parse_channels_only] Memory after gc.collect(): {process.memory_info().rss / 1024 / 1024:.2f} MB") - - if len(epgs_to_update) >= batch_size: - logger.info(f"[parse_channels_only] Bulk updating {len(epgs_to_update)} EPG entries") - if process: - logger.info(f"[parse_channels_only] Memory before bulk_update: {process.memory_info().rss / 1024 / 1024:.2f} MB") - EPGData.objects.bulk_update(epgs_to_update, ["name"]) - if process: - logger.info(f"[parse_channels_only] Memory after bulk_update: {process.memory_info().rss / 1024 / 1024:.2f} MB") - epgs_to_update = [] - # Force garbage collection - gc.collect() - - # Periodically clear the existing_epgs cache to prevent memory buildup - if processed_channels % 1000 == 0: - logger.info(f"[parse_channels_only] Clearing existing_epgs cache at {processed_channels} channels") - existing_epgs.clear() - gc.collect() - if process: - logger.info(f"[parse_channels_only] Memory after clearing cache: {process.memory_info().rss / 1024 / 1024:.2f} MB") - - # Send progress updates - if processed_channels % 100 == 0 or processed_channels == total_channels: - progress = 25 + int((processed_channels / total_channels) * 65) if total_channels > 0 else 90 - send_epg_update( - source.id, - "parsing_channels", - progress, - processed=processed_channels, - total=total_channels - ) - logger.debug(f"[parse_channels_only] Processed channel: {tvg_id} - {display_name}") - if process: - logger.info(f"[parse_channels_only] Memory before elem cleanup: {process.memory_info().rss / 1024 / 1024:.2f} MB") - # Clear memory - try: - # First clear the element's content + # If we encounter a programme element, we've processed all channels + # Break out of the loop to avoid memory spike + if elem.tag == 'programme': + logger.debug(f"[parse_channels_only] Found first programme element after processing {processed_channels} channels - exiting channel parsing") + # Clean up the element before breaking elem.clear() - - # Get the parent before we might lose reference to it parent = elem.getparent() if parent is not None: - # Clean up preceding siblings - while elem.getprevious() is not None: - del parent[0] + parent.remove(elem) + break - # Try to fully detach this element from parent - try: - parent.remove(elem) - del elem - del parent - except (ValueError, KeyError, TypeError): - # Element might already be removed or detached - pass - cleanup_memory(log_usage=True, force_collection=True) - time.sleep(.1) + # Only process channel elements + if elem.tag == 'channel': + channel_count += 1 + tvg_id = elem.get('id', '').strip() + if tvg_id: + display_name = None + for child in elem: + if child.tag == 'display-name' and child.text: + display_name = child.text.strip() + break - except Exception as e: - # Just log the error and continue - don't let cleanup errors stop processing - logger.debug(f"[parse_channels_only] Non-critical error during XML element cleanup: {e}") - if process: - logger.info(f"[parse_channels_only] Memory after elem cleanup: {process.memory_info().rss / 1024 / 1024:.2f} MB") - # Check if we should break early to avoid excessive sleep - if processed_channels >= total_channels and total_channels > 0: - logger.info(f"[parse_channels_only] Expected channel numbers hit, continuing - processed {processed_channels}/{total_channels}") - logger.debug(f"[parse_channels_only] Memory usage after {processed_channels}: {process.memory_info().rss / 1024 / 1024:.2f} MB") - #break - logger.debug(f"[parse_channels_only] Total elements processed: {total_elements_processed}") - # Add periodic forced cleanup based on TOTAL ELEMENTS, not just channels - # This ensures we clean up even if processing many non-channel elements - if total_elements_processed % 1000 == 0: - logger.info(f"[parse_channels_only] Performing preventative memory cleanup after {total_elements_processed} elements (found {processed_channels} channels)") - # Close and reopen the parser to release memory - if source_file and channel_parser: - # First clear element references - elem.clear() - if elem.getparent() is not None: - elem.getparent().remove(elem) + if not display_name: + display_name = tvg_id - # Reset parser state - del channel_parser - channel_parser = None + # Use lazy loading approach to reduce memory usage + if tvg_id in existing_tvg_ids: + # Only fetch the object if we need to update it and it hasn't been loaded yet + if tvg_id not in existing_epgs: + try: + # This loads the full EPG object from the database and caches it + existing_epgs[tvg_id] = EPGData.objects.get(tvg_id=tvg_id, epg_source=source) + except EPGData.DoesNotExist: + # Handle race condition where record was deleted + existing_tvg_ids.remove(tvg_id) + epgs_to_create.append(EPGData( + tvg_id=tvg_id, + name=display_name, + epg_source=source, + )) + logger.debug(f"[parse_channels_only] Added new channel to epgs_to_create 1: {tvg_id} - {display_name}") + continue + + # We use the cached object to check if the name has changed + epg_obj = existing_epgs[tvg_id] + if epg_obj.name != display_name: + # Only update if the name actually changed + epg_obj.name = display_name + epgs_to_update.append(epg_obj) + logger.debug(f"[parse_channels_only] Added channel to update to epgs_to_update: {tvg_id} - {display_name}") + else: + # This is a new channel that doesn't exist in our database + epgs_to_create.append(EPGData( + tvg_id=tvg_id, + name=display_name, + epg_source=source, + )) + logger.debug(f"[parse_channels_only] Added new channel to epgs_to_create 2: {tvg_id} - {display_name}") + + processed_channels += 1 + + # Batch processing + if len(epgs_to_create) >= batch_size: + logger.info(f"[parse_channels_only] Bulk creating {len(epgs_to_create)} EPG entries") + EPGData.objects.bulk_create(epgs_to_create, ignore_conflicts=True) + if process: + logger.info(f"[parse_channels_only] Memory after bulk_create: {process.memory_info().rss / 1024 / 1024:.2f} MB") + del epgs_to_create # Explicit deletion + epgs_to_create = [] + gc.collect() + if process: + logger.info(f"[parse_channels_only] Memory after gc.collect(): {process.memory_info().rss / 1024 / 1024:.2f} MB") + + if len(epgs_to_update) >= batch_size: + logger.info(f"[parse_channels_only] Bulk updating {len(epgs_to_update)} EPG entries") + if process: + logger.info(f"[parse_channels_only] Memory before bulk_update: {process.memory_info().rss / 1024 / 1024:.2f} MB") + EPGData.objects.bulk_update(epgs_to_update, ["name"]) + if process: + logger.info(f"[parse_channels_only] Memory after bulk_update: {process.memory_info().rss / 1024 / 1024:.2f} MB") + epgs_to_update = [] + # Force garbage collection gc.collect() - # Perform thorough cleanup - cleanup_memory(log_usage=True, force_collection=True) - - # Create a new parser context - channel_parser = etree.iterparse(source_file, events=('end',), tag='channel') - logger.info(f"[parse_channels_only] Recreated parser context after memory cleanup") - - # Also do cleanup based on processed channels as before - elif processed_channels % 1000 == 0 and processed_channels > 0: - logger.info(f"[parse_channels_only] Performing preventative memory cleanup at {processed_channels} channels") - # Close and reopen the parser to release memory - if source_file and channel_parser: - # First clear element references - elem.clear() - if elem.getparent() is not None: - elem.getparent().remove(elem) - - # Reset parser state - del channel_parser - channel_parser = None + # Periodically clear the existing_epgs cache to prevent memory buildup + if processed_channels % 1000 == 0: + logger.info(f"[parse_channels_only] Clearing existing_epgs cache at {processed_channels} channels") + existing_epgs.clear() gc.collect() + if process: + logger.info(f"[parse_channels_only] Memory after clearing cache: {process.memory_info().rss / 1024 / 1024:.2f} MB") - # Perform thorough cleanup + # Send progress updates + if processed_channels % 100 == 0 or processed_channels == total_channels: + progress = 25 + int((processed_channels / total_channels) * 65) if total_channels > 0 else 90 + send_epg_update( + source.id, + "parsing_channels", + progress, + processed=processed_channels, + total=total_channels + ) + logger.debug(f"[parse_channels_only] Processed channel: {tvg_id} - {display_name}") + if process: + logger.info(f"[parse_channels_only] Memory before elem cleanup: {process.memory_info().rss / 1024 / 1024:.2f} MB") + # Clear memory + try: + # First clear the element's content + elem.clear() + + # Get the parent before we might lose reference to it + parent = elem.getparent() + if parent is not None: + # Clean up preceding siblings + while elem.getprevious() is not None: + del parent[0] + + # Try to fully detach this element from parent + try: + parent.remove(elem) + del elem + del parent + except (ValueError, KeyError, TypeError): + # Element might already be removed or detached + pass cleanup_memory(log_usage=True, force_collection=True) + #time.sleep(.1) - # Create a new parser context - channel_parser = etree.iterparse(source_file, events=('end',), tag='channel') - logger.info(f"[parse_channels_only] Recreated parser context after memory cleanup") + except Exception as e: + # Just log the error and continue - don't let cleanup errors stop processing + logger.debug(f"[parse_channels_only] Non-critical error during XML element cleanup: {e}") + if process: + logger.info(f"[parse_channels_only] Memory after elem cleanup: {process.memory_info().rss / 1024 / 1024:.2f} MB") + # Check if we should break early to avoid excessive sleep + if processed_channels >= total_channels and total_channels > 0: + logger.info(f"[parse_channels_only] Expected channel numbers hit, continuing - processed {processed_channels}/{total_channels}") + logger.debug(f"[parse_channels_only] Memory usage after {processed_channels}: {process.memory_info().rss / 1024 / 1024:.2f} MB") + #break + logger.debug(f"[parse_channels_only] Total elements processed: {total_elements_processed}") + # Add periodic forced cleanup based on TOTAL ELEMENTS, not just channels + # This ensures we clean up even if processing many non-channel elements + if total_elements_processed % 1000 == 0: + logger.info(f"[parse_channels_only] Performing preventative memory cleanup after {total_elements_processed} elements (found {processed_channels} channels)") + # Close and reopen the parser to release memory + if source_file and channel_parser: + # First clear element references + elem.clear() + if elem.getparent() is not None: + elem.getparent().remove(elem) - if processed_channels == total_channels: - logger.info(f"[parse_channels_only] Processed all channels current memory: {process.memory_info().rss / 1024 / 1024:.2f} MB") + # Reset parser state + del channel_parser + channel_parser = None + gc.collect() + # Perform thorough cleanup + cleanup_memory(log_usage=True, force_collection=True) - if process: - logger.info(f"[parse_channels_only] Memory after leaving for loop: {process.memory_info().rss / 1024 / 1024:.2f} MB") - time.sleep(20) - # Explicit cleanup before sleeping - logger.info(f"[parse_channels_only] Completed channel parsing loop, processed {processed_channels} channels") - if process: - logger.info(f"[parse_channels_only] Memory before cleanup: {process.memory_info().rss / 1024 / 1024:.2f} MB") + # Create a new parser context - continue looking for both tags + # This doesn't restart from the beginning but continues from current position + channel_parser = etree.iterparse(source_file, events=('end',), tag=('channel', 'programme')) + logger.info(f"[parse_channels_only] Recreated parser context after memory cleanup") - # Explicit cleanup of the parser + # Also do cleanup based on processed channels as before + elif processed_channels % 1000 == 0 and processed_channels > 0: + logger.info(f"[parse_channels_only] Performing preventative memory cleanup at {processed_channels} channels") + # Close and reopen the parser to release memory + if source_file and channel_parser: + # First clear element references + elem.clear() + if elem.getparent() is not None: + elem.getparent().remove(elem) - #del channel_parser + # Reset parser state + del channel_parser + channel_parser = None + gc.collect() - logger.info(f"[parse_channels_only] Deleted channel_parser object") + # Perform thorough cleanup + cleanup_memory(log_usage=True, force_collection=True) - # Close the file - logger.info(f"[parse_channels_only] Closing file: {file_path}") - source_file.close() - logger.info(f"[parse_channels_only] File closed: {file_path}") + # Create a new parser context + # This doesn't restart from the beginning but continues from current position + channel_parser = etree.iterparse(source_file, events=('end',), tag=('channel', 'programme')) + logger.info(f"[parse_channels_only] Recreated parser context after memory cleanup") - # Force garbage collection - gc.collect() - if process: - logger.info(f"[parse_channels_only] Memory after final cleanup: {process.memory_info().rss / 1024 / 1024:.2f} MB") - - # Remove long sleep that might be causing issues - # time.sleep(200) # This seems excessive and may be causing issues + if processed_channels == total_channels: + logger.info(f"[parse_channels_only] Processed all channels current memory: {process.memory_info().rss / 1024 / 1024:.2f} MB") except (etree.XMLSyntaxError, Exception) as xml_error: logger.error(f"[parse_channels_only] XML parsing failed: {xml_error}") @@ -830,25 +828,18 @@ def parse_channels_only(source): send_epg_update(source.id, "parsing_channels", 100, status="error", error=str(xml_error)) return False if process: - logger.info(f"[parse_channels_only] Memory before final batch creation: {process.memory_info().rss / 1024 / 1024:.2f} MB") + logger.debug(f"[parse_channels_only] Memory before final batch creation: {process.memory_info().rss / 1024 / 1024:.2f} MB") # Process any remaining items if epgs_to_create: EPGData.objects.bulk_create(epgs_to_create, ignore_conflicts=True) - logger.info(f"[parse_channels_only] Created final batch of {len(epgs_to_create)} EPG entries") + logger.debug(f"[parse_channels_only] Created final batch of {len(epgs_to_create)} EPG entries") if epgs_to_update: EPGData.objects.bulk_update(epgs_to_update, ["name"]) - logger.info(f"[parse_channels_only] Updated final batch of {len(epgs_to_update)} EPG entries") + logger.debug(f"[parse_channels_only] Updated final batch of {len(epgs_to_update)} EPG entries") if process: - logger.info(f"[parse_channels_only] Memory after final batch creation: {process.memory_info().rss / 1024 / 1024:.2f} MB") - - # Final garbage collection and memory tracking - if process: - logger.info(f"[parse_channels_only] Memory before final gc: {process.memory_info().rss / 1024 / 1024:.2f} MB") - gc.collect() - if process: - logger.info(f"[parse_channels_only] Memory after final gc: {process.memory_info().rss / 1024 / 1024:.2f} MB") + logger.debug(f"[parse_channels_only] Memory after final batch creation: {process.memory_info().rss / 1024 / 1024:.2f} MB") # Update source status with channel count source.status = 'success' @@ -889,19 +880,34 @@ def parse_channels_only(source): return False finally: # Add more detailed cleanup in finally block - logger.info("In finally block, ensuring cleanup") - existing_tvg_ids = None - existing_epgs = None - - # Check final memory usage after clearing process - gc.collect() - # Add comprehensive cleanup at end of channel parsing - cleanup_memory(log_usage=True, force_collection=True) + logger.debug("In finally block, ensuring cleanup") try: - process = psutil.Process() - final_memory = process.memory_info().rss / 1024 / 1024 - logger.info(f"[parse_channels_only] Final memory usage: {final_memory:.2f} MB") - process = None + if 'channel_parser' in locals(): + del channel_parser + if 'elem' in locals(): + del elem + if 'parent' in locals(): + del parent + + if 'source_file' in locals(): + source_file.close() + del source_file + # Clear remaining large data structures + existing_epgs.clear() + epgs_to_create.clear() + epgs_to_update.clear() + existing_epgs = None + epgs_to_create = None + epgs_to_update = None + cleanup_memory(log_usage=True, force_collection=True) + except Exception as e: + logger.warning(f"Cleanup error: {e}") + + try: + if process: + final_memory = process.memory_info().rss / 1024 / 1024 + logger.debug(f"[parse_channels_only] Final memory usage: {final_memory:.2f} MB") + process = None except: pass