o
    `j                     @   s   d dl Z d dlZd dlZd dlZd dlZd dlZd dlmZ G dd dZG dd dej	Z
dd Zed	krBe  Zee dS dS )
    N)KafkaConsumerc                   @   s   e Zd Zedd ZdS )ConsumerPerformancec              	   C   s  zi }| j D ]1}|d\}}zt|}W n	 ty   Y nw |dkr&d }n|dkr-d}n|dkr3d}|||< qtd | j|d< d	|d
< d|vrMd|d< | jd |d< | D ]\}}td|| qXt	| j
fi |}td| j
 td| j td| j t }t| j||| jd}|  td t  t }d}	|D ]}
|	d7 }	|	| jkr nqt }|  |  td|	 td|| d W d S  ty   t }tj|  td Y d S w )N=NoneFalseFTrueTzInitializing Consumer...bootstrap_serversearliestauto_offset_resetconsumer_timeout_msi'  i  metrics_sample_window_msz---> {0}={1}z---> topic={0}z ---> report stats every {0} secsz---> raw metrics? {0})eventraw_metricsz-> OK!r      zConsumed {0} recordszExecution time:secs)consumer_configsplitint
ValueErrorprintr   stats_intervalitemsformatr   topicr   	threadingEventStatsReporterstarttime	monotonicnum_recordssetjoin	Exceptionsysexc_info	tracebackprint_exceptionexit)argspropspropkvconsumer
timer_stoptimer
start_timerecordsmsgend_timer%    r5   c/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/benchmarks/consumer_performance.pyrun   sj   




zConsumerPerformance.runN)__name__
__module____qualname__staticmethodr7   r5   r5   r5   r6   r      s    r   c                       s6   e Zd Zd fdd	Zdd Zdd Zd	d
 Z  ZS )r   NFc                    s&   t    || _|| _|| _|| _d S N)super__init__intervalr.   r   r   )selfr?   r.   r   r   	__class__r5   r6   r>   I   s
   

zStatsReporter.__init__c                 C   s:   | j  }| jrt| d S tdjdi |d  d S )Na  {records-consumed-rate:.0f} records/sec ({bytes-consumed-rate:.0f} B/sec), {fetch-latency-avg:.0f}ms avg latency, {fetch-rate:.0f} avg fetch requests/sec, {fetch-size-avg:.0f} avg fetch size, {records-lag-max:.0f} max record lag, {records-per-request-avg:.0f} avg records/reqzconsumer-fetch-manager-metricsr5   )r.   metricsr   pprintr   r   )r@   rC   r5   r5   r6   print_statsP   s   
zStatsReporter.print_statsc                 C   s   |    d S r<   )rE   r@   r5   r5   r6   print_final^   s   zStatsReporter.print_finalc                 C   s<   | j r| j | js|   | j r| j | jr
|   d S r<   )r   waitr?   rE   rG   rF   r5   r5   r6   r7   a   s   zStatsReporter.run)NF)r8   r9   r:   r>   rE   rG   r7   __classcell__r5   r5   rA   r6   r   H   s
    r   c                  C   s   t jdd} | jddtdddd | jd	d
tddd | jdtddd | jddtdddd | jdtdd | jdtddd | jdddd | S )Nz5This tool is used to verify the consumer performance.)descriptionz-bz--bootstrap-servers+r5   z'host:port for cluster bootstrap servers)typenargsdefaulthelpz-tz--topicz>Topic for consumer test (default: kafka-python-benchmark-test)zkafka-python-benchmark-test)rL   rO   rN   z--num-recordsz0number of messages to consume (default: 1000000)i@B z-cz--consumer-configzVkafka consumer related configuration properties like bootstrap_servers,client_id etc..z--fixture-compressionzBspecify a compression type for use with broker fixtures / producer)rL   rO   z--stats-intervalz?Interval in seconds for stats reporting to console (default: 5)   z--raw-metrics
store_truez<Enable this flag to print full metrics dict on each interval)actionrO   )argparseArgumentParseradd_argumentstrr   )parserr5   r5   r6   get_args_parserh   sF   

rX   __main__)rS   rD   r$   r   r   r&   kafkar   r   Threadr   rX   r8   
parse_argsr)   r7   r5   r5   r5   r6   <module>   s   :  
