Téléchargement du journal PostgreSQL depuis le cloud AWS

Ou un peu de tetris appliqué.
Tout nouveau est bien oublié ancien.
Épigraphes.
Téléchargement du journal PostgreSQL depuis le cloud AWS

Définition du problÚme

Il est nécessaire de télécharger périodiquement le fichier de log actuel de PostgreSQL depuis le cloud AWS sur un hÎte Linux local. Pas en temps réel, mais disons avec un léger retard.
La période de chargement du fichier de log est de 5 minutes.
Le fichier de log, dans AWS, est tourné chaque heure.

Outils utilisés

Pour télécharger le fichier de log sur l'hÎte, un script bash est utilisé, invoquant l'API AWS «aws rds download-db-log-file-portion».

ParamĂštres :

  • —db-instance-identifier : Le nom de l'instance dans AWS ;
  • —log-file-name : le nom du fichier de log actuel gĂ©nĂ©rĂ©
  • —max-item : Le nombre total d'Ă©lĂ©ments renvoyĂ©s dans les rĂ©sultats de la commande.Taille de la portion du fichier tĂ©lĂ©chargĂ©.
  • —starting-token : Jeton de la portion initiale

Dans ce cas particulier, la tĂąche de tĂ©lĂ©chargement des logs a Ă©tĂ© levĂ©e dans le cadre de la surveillance des performances des requĂȘtes PostgreSQL.

Et c'est tout simplement une tùche intéressante, pour s'entraßner et varier durant le temps de travail.
Je suppose que la tùche, de par sa banalité, a déjà été résolue. Mais un rapide google des solutions n'a rien donné, et il n'y avait pas vraiment d'envie de chercher plus en profondeur. Quoi qu'il en soit, c'est un bon entraßnement.

Formalisation de la tĂąche

Le fichier de log final se compose de nombreuses lignes de longueur variable. Graphiquement, le fichier de log peut ĂȘtre reprĂ©sentĂ© ainsi :
Téléchargement du journal PostgreSQL depuis le cloud AWS

Ça ressemble dĂ©jĂ  Ă  quelque chose ? Quel rapport avec « tetris » ? Eh bien, justement.
Si l'on reprĂ©sente graphiquement les variantes possibles qui se prĂ©sentent lors du chargement d'un nouveau fichier (pour simplifier, dans ce cas, disons que les lignes ont toutes la mĂȘme longueur), on obtient les formes standards de tetris :

1) Le fichier est entiÚrement chargé et est final. La taille de la portion est supérieure à celle du fichier final :
Téléchargement du journal PostgreSQL depuis le cloud AWS

2) Le fichier a une continuation. La taille de la portion est inférieure à celle du fichier final :
Téléchargement du journal PostgreSQL depuis le cloud AWS

3) Le fichier est une continuation du fichier précédent et a une continuation. La taille de la portion est inférieure au reste du fichier final :
Téléchargement du journal PostgreSQL depuis le cloud AWS

4) Le fichier est une continuation du fichier précédent et est final. La taille de la portion est supérieure au reste du fichier final :
Téléchargement du journal PostgreSQL depuis le cloud AWS

La tĂąche consiste Ă  assembler un rectangle ou Ă  jouer au tetris, Ă  un nouveau niveau.
Téléchargement du journal PostgreSQL depuis le cloud AWS

ProblÚmes rencontrés lors de la résolution de la tùche

1) Assembler une chaĂźne Ă  partir de 2 portions

Téléchargement du journal PostgreSQL depuis le cloud AWS
En fait, aucun problĂšme particulier n'a surgi. TĂąche standard d'un cours de programmation d'introduction.

Taille optimale de la portion

Et cela, c'est un peu plus intéressant.
Malheureusement, il n'est pas possible d'utiliser un décalage aprÚs le jeton de la premiÚre portion :

Comme vous le savez dĂ©jĂ , l'option —starting-token est utilisĂ©e pour spĂ©cifier le point de dĂ©part de la pagination. Cette option prend des valeurs de chaĂźne, ce qui signifie que si vous essayez d'ajouter une valeur de dĂ©calage devant la chaĂźne du jeton suivant, l'option ne sera pas considĂ©rĂ©e comme un dĂ©calage.

C'est pourquoi il faut lire par morceaux.
Si vous lisez en grandes portions, le nombre de lectures sera minimal, mais le volume sera maximal.
Si vous lisez en petites portions, au contraire, le nombre de lectures sera maximal, mais le volume sera minimal.
Par conséquent, pour réduire le trafic et améliorer l'esthétique de la solution, il a fallu concevoir un certain type de solution qui ressemble malheureusement à un palliatif.

Pour illustrer, considérons le processus de chargement d'un fichier journal dans 2 variantes fortement simplifiées. Le nombre de lectures dans les deux cas dépend de la taille de la portion.

1) Chargement par petites portions :
Téléchargement du journal PostgreSQL depuis le cloud AWS

2) Chargement par grandes portions :
Téléchargement du journal PostgreSQL depuis le cloud AWS

Comme d'habitude, la solution optimale se trouve au milieu..
La taille de la portion est minimale, mais au cours de la lecture, la taille peut ĂȘtre augmentĂ©e pour rĂ©duire le nombre de lectures.

Il convient de noter que la tĂąche de dĂ©terminer la taille optimale de la portion lue n'est pas encore complĂštement rĂ©solue et nĂ©cessite une analyse plus approfondie. Peut-ĂȘtre un peu plus tard.

Description générale de l'implémentation

Tables de service utilisées

CREATE TABLE endpoint
(
id SERIAL ,
host text 
);

TABLE database
(
id SERIAL , 


last_aws_log_time text ,
last_aws_nexttoken text ,
aws_max_item_size integer 
);
last_aws_log_time — horodatage du dernier fichier journal chargĂ© au format YYYY-MM-DD-HH24.
last_aws_nexttoken — Ă©tiquette textuelle de la derniĂšre portion chargĂ©e.
aws_max_item_size - taille de portion initiale déterminée empiriquement.

Texte complet du script

download_aws_piece.sh

#!/bin/bash
#########################################################
# download_aws_piece.sh
# downloan piece of log from AWS
# version HABR
 let min_item_size=1024
 let max_item_size=1048576
 let growth_factor=3
 let growth_counter=1
 let growth_counter_max=3

 echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:''STARTED'
 
 AWS_LOG_TIME=$1
 echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:AWS_LOG_TIME='$AWS_LOG_TIME
  
 database_id=$2
 echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:database_id='$database_id
 RESULT_FILE=$3 
  
 endpoint=`psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE_DATABASE -A -t -c "select e.host from endpoint e join database d on e.id = d.endpoint_id where d.id = $database_id "`
 echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:endpoint='$endpoint
  
 db_instance=`echo $endpoint | awk -F"." '{print toupper($1)}'`
 
 echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:db_instance='$db_instance

 LOG_FILE=$RESULT_FILE'.tmp_log'
 TMP_FILE=$LOG_FILE'.tmp'
 TMP_MIDDLE=$LOG_FILE'.tmp_mid'  
 TMP_MIDDLE2=$LOG_FILE'.tmp_mid2'  
  
 current_aws_log_time=`psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -A -t -c "select last_aws_log_time from database where id = $database_id "`

 echo $(date +%Y%m%d%H%M)':      download_aws_piece.sh:current_aws_log_time='$current_aws_log_time
  
  if [[ $current_aws_log_time != $AWS_LOG_TIME  ]];
  then
    is_new_log='1'
	if ! psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -v ON_ERROR_STOP=1 -A -t -q -c "update database set last_aws_log_time = '$AWS_LOG_TIME' where id = $database_id "
	then
	  echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: FATAL_ERROR - update database set last_aws_log_time .'
	  exit 1
	fi
  else
    is_new_log='0'
  fi
  
  echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:is_new_log='$is_new_log
  
  let last_aws_max_item_size=`psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -A -t -c "select aws_max_item_size from database where id = $database_id "`
  echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: last_aws_max_item_size='$last_aws_max_item_size
  
  let count=1
  if [[ $is_new_log == '1' ]];
  then    
	echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: START DOWNLOADING OF NEW AWS LOG'
	if ! aws rds download-db-log-file-portion 
		--max-items $last_aws_max_item_size 
		--region REGION 
		--db-instance-identifier  $db_instance 
		--log-file-name error/postgresql.log.$AWS_LOG_TIME > $LOG_FILE
	then
		echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: FATAL_ERROR - Could not get log from AWS .'
		exit 2
	fi  	
  else
    next_token=`psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -v ON_ERROR_STOP=1 -A -t -c "select last_aws_nexttoken from database where id = $database_id "`
	
	if [[ $next_token == '' ]];
	then
	  next_token='0'	  
	fi
	
	echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: CONTINUE DOWNLOADING OF AWS LOG'
	if ! aws rds download-db-log-file-portion 
	    --max-items $last_aws_max_item_size 
		--starting-token $next_token 
		--region REGION 
		--db-instance-identifier  $db_instance 
		--log-file-name error/postgresql.log.$AWS_LOG_TIME > $LOG_FILE
	then
		echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: FATAL_ERROR - Could not get log from AWS .'
		exit 3
	fi       
	
	line_count=`cat  $LOG_FILE | wc -l`
	let lines=$line_count-1
	  
	tail -$lines $LOG_FILE > $TMP_MIDDLE 
	mv -f $TMP_MIDDLE $LOG_FILE
  fi
  
  next_token_str=`cat $LOG_FILE | grep NEXTTOKEN` 
  next_token=`echo $next_token_str | awk -F" " '{ print $2}' `
  
  grep -v NEXTTOKEN $LOG_FILE  > $TMP_FILE 
  
  if [[ $next_token == '' ]];
  then
	  cp $TMP_FILE $RESULT_FILE
	  
	  echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:  NEXTTOKEN NOT FOUND - FINISH '
	  rm $LOG_FILE 
	  rm $TMP_FILE
	  rm $TMP_MIDDLE
          rm $TMP_MIDDLE2	  
	  exit 0  
  else
	psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -v ON_ERROR_STOP=1 -A -t -q -c "update database set last_aws_nexttoken = '$next_token' where id = $database_id "
  fi
  
  first_str=`tail -1 $TMP_FILE`
  
  line_count=`cat  $TMP_FILE | wc -l`
  let lines=$line_count-1    
  
  head -$lines $TMP_FILE  > $RESULT_FILE

###############################################
# MAIN CIRCLE
  let count=2
  while [[ $next_token != '' ]];
  do 
    echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: count='$count
	
	echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: START DOWNLOADING OF AWS LOG'
	if ! aws rds download-db-log-file-portion 
             --max-items $last_aws_max_item_size 
             --starting-token $next_token 
             --region REGION 
             --db-instance-identifier  $db_instance 
             --log-file-name error/postgresql.log.$AWS_LOG_TIME > $LOG_FILE
	then
		echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: FATAL_ERROR - Could not get log from AWS .'
		exit 4
	fi

	next_token_str=`cat $LOG_FILE | grep NEXTTOKEN` 
	next_token=`echo $next_token_str | awk -F" " '{ print $2}' `

	TMP_FILE=$LOG_FILE'.tmp'
	grep -v NEXTTOKEN $LOG_FILE  > $TMP_FILE  
	
	last_str=`head -1 $TMP_FILE`
  
    if [[ $next_token == '' ]];
	then
	  concat_str=$first_str$last_str
	  	  
	  echo $concat_str >> $RESULT_FILE
		 
	  line_count=`cat  $TMP_FILE | wc -l`
	  let lines=$line_count-1
	  
	  tail -$lines $TMP_FILE >> $RESULT_FILE
	  
	  echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:  NEXTTOKEN NOT FOUND - FINISH '
	  rm $LOG_FILE 
	  rm $TMP_FILE
	  rm $TMP_MIDDLE
          rm $TMP_MIDDLE2	  
	  exit 0  
	fi
	
    if [[ $next_token != '' ]];
	then
		let growth_counter=$growth_counter+1
		if [[ $growth_counter -gt $growth_counter_max ]];
		then
			let last_aws_max_item_size=$last_aws_max_item_size*$growth_factor
			let growth_counter=1
		fi
	
		if [[ $last_aws_max_item_size -gt $max_item_size ]]; 
		then
			let last_aws_max_item_size=$max_item_size
		fi 

	  psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -A -t -q -c "update database set last_aws_nexttoken = '$next_token' where id = $database_id "
	  
	  concat_str=$first_str$last_str
	  	  
	  echo $concat_str >> $RESULT_FILE
		 
	  line_count=`cat  $TMP_FILE | wc -l`
	  let lines=$line_count-1
	  
	  #############################
	  #Get middle of file
	  head -$lines $TMP_FILE > $TMP_MIDDLE
	  
	  line_count=`cat  $TMP_MIDDLE | wc -l`
	  let lines=$line_count-1
	  tail -$lines $TMP_MIDDLE > $TMP_MIDDLE2
	  
	  cat $TMP_MIDDLE2 >> $RESULT_FILE	  
	  
	  first_str=`tail -1 $TMP_FILE`	  
	fi
	  
    let count=$count+1

  done
#
#################################################################

exit 0  

Fragments du script avec quelques explications :

ParamÚtres d'entrée du script :

  • Horodatage du nom du fichier journal au format YYYY-MM-DD-HH24 : AWS_LOG_TIME=$1
  • ID de la base de donnĂ©es : database_id=$2
  • Nom du fichier journal assemblĂ© : RESULT_FILE=$3

Obtenir l'horodatage du dernier fichier journal chargé :

current_aws_log_time=`psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -A -t -c "select last_aws_log_time from database where id = $database_id "`

Si l'horodatage du dernier fichier journal chargé ne correspond pas au paramÚtre d'entrée, un nouveau fichier journal sera téléchargé :

si [[ $current_aws_log_time != $AWS_LOG_TIME  ]];
  alors
    is_new_log='1'
	si ! psql -h ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -v ON_ERROR_STOP=1 -A -t -c "update database set last_aws_log_time = '$AWS_LOG_TIME' where id = $database_id "
	alors
	  echo '***download_aws_piece.sh -ERREUR_FATALE - mise à jour de la base de données pour définir last_aws_log_time.'
	  exit 1
	fi
  sinon
    is_new_log='0'
  fi

Nous récupérons la valeur de l'étiquette nexttoken à partir du fichier téléchargé :

  next_token_str=`cat $LOG_FILE | grep NEXTTOKEN` 
  next_token=`echo $next_token_str | awk -F" " '{ print $2}' `

Le critĂšre de fin de chargement est la valeur vide de nexttoken.

Dans la boucle, nous comptons les portions du fichier, tout en enchaĂźnant les lignes et en augmentant la taille de la portion :
Boucle principale

# MAIN CIRCLE
  let count=2
  while [[ $next_token != '' ]];
  do 
    echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: count='$count
	
	echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: START DOWNLOADING OF AWS LOG'
	if ! aws rds download-db-log-file-portion 
     --max-items $last_aws_max_item_size 
	 --starting-token $next_token 
     --region REGION 
     --db-instance-identifier  $db_instance 
     --log-file-name error/postgresql.log.$AWS_LOG_TIME > $LOG_FILE
	then
		echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh: FATAL_ERROR - Could not get log from AWS .'
		exit 4
	fi

	next_token_str=`cat $LOG_FILE | grep NEXTTOKEN` 
	next_token=`echo $next_token_str | awk -F" " '{ print $2}' `

	TMP_FILE=$LOG_FILE'.tmp'
	grep -v NEXTTOKEN $LOG_FILE  > $TMP_FILE  
	
	last_str=`head -1 $TMP_FILE`
  
    if [[ $next_token == '' ]];
	then
	  concat_str=$first_str$last_str
	  	  
	  echo $concat_str >> $RESULT_FILE
		 
	  line_count=`cat  $TMP_FILE | wc -l`
	  let lines=$line_count-1
	  
	  tail -$lines $TMP_FILE >> $RESULT_FILE
	  
	  echo $(date +%Y%m%d%H%M)':    download_aws_piece.sh:  NEXTTOKEN NOT FOUND - FINISH '
	  rm $LOG_FILE 
	  rm $TMP_FILE
	  rm $TMP_MIDDLE
         rm $TMP_MIDDLE2	  
	  exit 0  
	fi
	
    if [[ $next_token != '' ]];
	then
		let growth_counter=$growth_counter+1
		if [[ $growth_counter -gt $growth_counter_max ]];
		then
			let last_aws_max_item_size=$last_aws_max_item_size*$growth_factor
			let growth_counter=1
		fi
	
		if [[ $last_aws_max_item_size -gt $max_item_size ]]; 
		then
			let last_aws_max_item_size=$max_item_size
		fi 

	  psql -h MONITOR_ENDPOINT.rds.amazonaws.com -U USER -d MONITOR_DATABASE -A -t -q -c "update database set last_aws_nexttoken = '$next_token' where id = $database_id "
	  
	  concat_str=$first_str$last_str
	  	  
	  echo $concat_str >> $RESULT_FILE
		 
	  line_count=`cat  $TMP_FILE | wc -l`
	  let lines=$line_count-1
	  
	  #############################
	  #Get middle of file
	  head -$lines $TMP_FILE > $TMP_MIDDLE
	  
	  line_count=`cat  $TMP_MIDDLE | wc -l`
	  let lines=$line_count-1
	  tail -$lines $TMP_MIDDLE > $TMP_MIDDLE2
	  
	  cat $TMP_MIDDLE2 >> $RESULT_FILE	  
	  
	  first_str=`tail -1 $TMP_FILE`	  
	fi
	  
    let count=$count+1

  done

Et aprĂšs ?

Ainsi, la premiĂšre tĂąche intermĂ©diaire – « tĂ©lĂ©charger le fichier log depuis le cloud » est rĂ©solue. Que faire avec le log tĂ©lĂ©chargĂ© ?
Pour commencer, il est nĂ©cessaire d'analyser le fichier log et d'en extraire les requĂȘtes.
La tĂąche n'est pas trĂšs difficile. Un simple script bash suffit.
upload_log_query.sh

#!/bin/bash
#########################################################
# upload_log_query.sh
# Upload table table from dowloaded aws file 
# version HABR
###########################################################  
echo 'TIMESTAMP:'$(date +%c)' Upload log_query table '
source_file=$1
echo 'source_file='$source_file
database_id=$2
echo 'database_id='$database_id

beginer=' '
first_line='1'
let "line_count=0"
sql_line=' '
sql_flag=' '    
space=' '
cat $source_file | while read line
do
  line="$space$line"

  if [[ $first_line == "1" ]]; then
    beginer=`echo $line | awk -F" " '{ print $1}' `
    first_line='0'
  fi

  current_beginer=`echo $line | awk -F" " '{ print $1}' `

  if [[ $current_beginer == $beginer ]]; then
    if [[ $sql_flag == '1' ]]; then
     sql_flag='0' 
     log_date=`echo $sql_line | awk -F" " '{ print $1}' `
     log_time=`echo $sql_line | awk -F" " '{ print $2}' `
     duration=`echo $sql_line | awk -F" " '{ print $5}' `

     #replace ' to ''
     sql_modline=`echo "$sql_line" | sed 's/'''/''''''/g'`
     sql_line=' '

	 ################
	 #PROCESSING OF THE SQL-SELECT IS HERE
     if ! psql -h ENDPOINT.rds.amazonaws.com -U USER -d DATABASE -v ON_ERROR_STOP=1 -A -t -c "select log_query('$ip_port',$database_id , '$log_date' , '$log_time' , '$duration' , '$sql_modline' )" 
     then
        echo 'FATAL_ERROR - log_query '
        exit 1
     fi
	 ################

    fi #if [[ $sql_flag == '1' ]]; then

    let "line_count=line_count+1"

    check=`echo $line | awk -F" " '{ print $8}' `
    check_sql=${check^^}    

    #echo 'check_sql='$check_sql
    
    if [[ $check_sql == 'SELECT' ]]; then
     sql_flag='1'    
     sql_line="$sql_line$line"
	 ip_port=`echo $sql_line | awk -F":" '{ print $4}' `
    fi
  else       

    if [[ $sql_flag == '1' ]]; then
      sql_line="$sql_line$line"
    fi   
    
  fi #if [[ $current_beginer == $beginer ]]; then

done

Maintenant, avec la requĂȘte extraite du fichier log, nous pouvons travailler.

Et plusieurs opportunités utiles s'offrent à nous.

Les requĂȘtes analysĂ©es doivent ĂȘtre stockĂ©es quelque part. Pour cela, nous utilisons une table de service log_query

CREATE TABLE log_query
(
   id SERIAL ,
   queryid bigint ,
   query_md5hash text not null ,
   database_id integer not null ,  
   timepoint timestamp without time zone not null,
   duration double precision not null ,
   query text not null ,
   explained_plan text[],
   plan_md5hash text  , 
   explained_plan_wo_costs text[],
   plan_hash_value text  ,
   baseline_id integer ,
   ip text ,
   port text 
);
ALTER TABLE log_query ADD PRIMARY KEY (id);
ALTER TABLE log_query ADD CONSTRAINT queryid_timepoint_unique_key UNIQUE (queryid, timepoint );
ALTER TABLE log_query ADD CONSTRAINT query_md5hash_timepoint_unique_key UNIQUE (query_md5hash, timepoint );

CREATE INDEX log_query_timepoint_idx ON log_query (timepoint);
CREATE INDEX log_query_queryid_idx ON log_query (queryid);
ALTER TABLE log_query ADD CONSTRAINT database_id_fk FOREIGN KEY (database_id) REFERENCES database (id) ON DELETE CASCADE ;

Le traitement de la requĂȘte analysĂ©e se fait dans plpgsql la fonction «log_query».
log_query.sql

--log_query.sql
--version HABR
CREATE OR REPLACE FUNCTION log_query( ip_port text ,log_database_id integer , log_date text , log_time text , duration text , sql_line text   ) RETURNS boolean AS $$
DECLARE
  result boolean ;
  log_timepoint timestamp without time zone ;
  log_duration double precision ; 
  pos integer ;
  log_query text ;
  activity_string text ;
  log_md5hash text ;
  log_explain_plan text[] ;
  
  log_planhash text ;
  log_plan_wo_costs text[] ; 
  
  database_rec record ;
  
  pg_stat_query text ; 
  test_log_query text ;
  log_query_rec record;
  found_flag boolean;
  
  pg_stat_history_rec record ;
  port_start integer ;
  port_end integer ;
  client_ip text ;
  client_port text ;
  log_queryid bigint ;
  log_query_text text ;
  pg_stat_query_text text ; 
BEGIN
  result = TRUE ;

  RAISE NOTICE '***log_query';
  
  port_start = position('(' in ip_port);
  port_end = position(')' in ip_port);
  client_ip = substring( ip_port from 1 for port_start-1 );
  client_port = substring( ip_port from port_start+1 for port_end-port_start-1 );

  SELECT e.host , d.name , d.owner_pwd 
  INTO database_rec
  FROM database d JOIN endpoint e ON e.id = d.endpoint_id
  WHERE d.id = log_database_id ;
  
  log_timepoint = to_timestamp(log_date||' '||log_time,'YYYY-MM-DD HH24-MI-SS');
  log_duration = duration:: double precision; 

  
  pos = position ('SELECT' in UPPER(sql_line) );
  log_query = substring( sql_line from pos for LENGTH(sql_line));
  log_query = regexp_replace(log_query,' +',' ','g');
  log_query = regexp_replace(log_query,';+','','g');
  log_query = trim(trailing ' ' from log_query);
 

  log_md5hash = md5( log_query::text );
  
  --Explain execution plan--
  EXECUTE 'SELECT dblink_connect(''LINK1'',''host='||database_rec.host||' dbname='||database_rec.name||' user=DATABASE password='||database_rec.owner_pwd||' '')'; 
  
  log_explain_plan = ARRAY ( SELECT * FROM dblink('LINK1', 'EXPLAIN '||log_query ) AS t (plan text) );
  log_plan_wo_costs = ARRAY ( SELECT * FROM dblink('LINK1', 'EXPLAIN ( COSTS FALSE ) '||log_query ) AS t (plan text) );
    
  PERFORM dblink_disconnect('LINK1');
  --------------------------
  BEGIN
	INSERT INTO log_query
	(
		query_md5hash ,
		database_id , 
		timepoint ,
		duration ,
		query ,
		explained_plan ,
		plan_md5hash , 
		explained_plan_wo_costs , 
		plan_hash_value , 
		ip , 
		port
	) 
	VALUES 
	(
		log_md5hash ,
		log_database_id , 
		log_timepoint , 
		log_duration , 
		log_query ,
		log_explain_plan , 
		md5(log_explain_plan::text) ,
		log_plan_wo_costs , 
		md5(log_plan_wo_costs::text),
		client_ip , 
		client_port		
	);
	activity_string = 	'Nouvelle requĂȘte enregistrĂ©e '||
						' database_id = '|| log_database_id ||
						' query_md5hash='||log_md5hash||
						' , timepoint = '||to_char(log_timepoint,'YYYYMMDD HH24:MI:SS');
					
	RAISE NOTICE '%',activity_string;					
					 
	PERFORM pg_log( log_database_id , 'log_query' , activity_string);  

	EXCEPTION
	  WHEN unique_violation THEN
		RAISE NOTICE '*** unique_violation *** la requĂȘte est dĂ©jĂ  enregistrĂ©e';
	END;

	SELECT 	queryid
	INTO   	log_queryid
	FROM 	log_query 
	WHERE 	query_md5hash = log_md5hash AND
			timepoint = log_timepoint;

	IF log_queryid IS NOT NULL 
	THEN 
	  RAISE NOTICE 'log_query avec query_md5hash = % et timepoint = % a déjà un QUERYID = %',log_md5hash,log_timepoint , log_queryid ;
	  RETURN result;
	END IF;
	
	------------------------------------------------
	RAISE NOTICE 'Mise Ă  jour du queryid';	
	
	SELECT * 
	INTO log_query_rec
	FROM log_query
	WHERE query_md5hash = log_md5hash AND timepoint = log_timepoint ; 
	
	log_query_rec.query=regexp_replace(log_query_rec.query,';+','','g');
	
	FOR pg_stat_history_rec IN
	 SELECT 
         queryid ,
	  query 
	 FROM 
         pg_stat_db_queries 
     WHERE  
      database_id = log_database_id AND
       queryid is not null 
	LOOP
	  pg_stat_query = pg_stat_history_rec.query ; 
	  pg_stat_query=regexp_replace(pg_stat_query,'n+',' ','g');
	  pg_stat_query=regexp_replace(pg_stat_query,'t+',' ','g');
	  pg_stat_query=regexp_replace(pg_stat_query,' +',' ','g');
	  pg_stat_query=regexp_replace(pg_stat_query,'$.','%','g');
	
	  log_query_text = trim(trailing ' ' from log_query_rec.query);
	  pg_stat_query_text = pg_stat_query; 
	
	  
	  --SELECT log_query_rec.query like pg_stat_query INTO found_flag ; 
	  IF (log_query_text LIKE pg_stat_query_text) THEN
		found_flag = TRUE ;
	  ELSE
		found_flag = FALSE ;
	  END IF;	  
	  
	  
	  IF found_flag THEN
	    
		UPDATE log_query SET queryid = pg_stat_history_rec.queryid WHERE query_md5hash = log_md5hash AND timepoint = log_timepoint ;
		activity_string = 	' updated queryid = '||pg_stat_history_rec.queryid||
		                    ' pour log_query avec id = '||log_query_rec.id               
		   				    ;						
	    RAISE NOTICE '%',activity_string;	
		EXIT ;
	  END IF ;
	  
	END LOOP ;
	
  RETURN result ;
END
$$ LANGUAGE plpgsql;

Lors du traitement, un tableau de service est utilisĂ© pg_stat_db_queries, contenant un instantanĂ© des requĂȘtes en cours Ă  partir du tableau pg_stat_history (L'utilisation du tableau est dĂ©crite ici — Surveillance des performances des requĂȘtes PostgreSQL. Partie 1 — rapport)

TABLE pg_stat_db_queries
(
   database_id integer,  
   queryid bigint ,  
   query text , 
   max_time double precision 
);

TABLE pg_stat_history 
(


database_id integer ,


queryid bigint ,


max_time double precision	 , 	


);

La fonction permet d'effectuer un certain nombre de fonctionnalitĂ©s utiles pour le traitement des requĂȘtes Ă  partir du fichier journal. En particulier :

PossibilitĂ© n°1 — Historique des exĂ©cutions de requĂȘtes

TrĂšs utile pour commencer Ă  rĂ©soudre un incident de performance. D'abord, consulter l'historique — quand le ralentissement a-t-il commencĂ© ?
Ensuite, selon la mĂ©thode classique – rechercher des causes externes. Peut-ĂȘtre que la charge de la base de donnĂ©es a simplement augmentĂ© de maniĂšre brusque et qu'une requĂȘte spĂ©cifique n'y est pour rien.
Ajouter un nouvel enregistrement dans le tableau log_query

  port_start = position('(' in ip_port);
  port_end = position(')' in ip_port);
  client_ip = substring( ip_port from 1 for port_start-1 );
  client_port = substring( ip_port from port_start+1 for port_end-port_start-1 );

  SELECT e.host , d.name , d.owner_pwd 
  INTO database_rec
  FROM database d JOIN endpoint e ON e.id = d.endpoint_id
  WHERE d.id = log_database_id ;
  
  log_timepoint = to_timestamp(log_date||' '||log_time,'YYYY-MM-DD HH24-MI-SS');
  log_duration = to_number(duration,'99999999999999999999D9999999999'); 

  
  pos = position ('SELECT' in UPPER(sql_line) );
  log_query = substring( sql_line from pos for LENGTH(sql_line));
  log_query = regexp_replace(log_query,' +',' ','g');
  log_query = regexp_replace(log_query,';+','','g');
  log_query = trim(trailing ' ' from log_query);
 
  RAISE NOTICE 'log_query=%',log_query ;   

  log_md5hash = md5( log_query::text );
  
  --Expliquer le plan d'exécution--
  EXECUTE 'SELECT dblink_connect(''LINK1'',''host='||database_rec.host||' dbname='||database_rec.name||' user=DATABASE password='||database_rec.owner_pwd||' '')'; 
  
  log_explain_plan = ARRAY ( SELECT * FROM dblink('LINK1', 'EXPLAIN '||log_query ) AS t (plan text) );
  log_plan_wo_costs = ARRAY ( SELECT * FROM dblink('LINK1', 'EXPLAIN ( COSTS FALSE ) '||log_query ) AS t (plan text) );
    
  PERFORM dblink_disconnect('LINK1');
  --------------------------
  BEGIN
	INSERT INTO log_query
	(
		query_md5hash ,
		database_id , 
		timepoint ,
		duration ,
		query ,
		explained_plan ,
		plan_md5hash , 
		explained_plan_wo_costs , 
		plan_hash_value , 
		ip , 
		port
	) 
	VALUES 
	(
		log_md5hash ,
		log_database_id , 
		log_timepoint ,
		log_duration ,
		log_query ,
		log_explain_plan , 
		md5(log_explain_plan::text) ,
		log_plan_wo_costs , 
		md5(log_plan_wo_costs::text),
		client_ip , 
		client_port		
	);

PossibilitĂ© n°2 — Sauvegarder les plans d'exĂ©cution des requĂȘtes

À ce stade, une objection-prĂ©cision-commentaire peut ĂȘtre soulevĂ©e : «Mais il y a dĂ©jĂ  autoexplain». C'est bien qu'il existe, mais quel en est l'intĂ©rĂȘt si le plan d'exĂ©cution est stockĂ© dans le mĂȘme fichier journal et que pour le conserver pour une analyse ultĂ©rieure, il faut alors parser le fichier journal ?

Pour moi, il me fallait :
premiÚrement : stocker le plan d'exécution dans un tableau de service de la base de données de surveillance ;
DeuxiĂšmement : avoir la possibilitĂ© de comparer les plans d'exĂ©cution entre eux, afin de voir immĂ©diatement si le plan d'exĂ©cution de la requĂȘte a changĂ©.

Une requĂȘte avec des paramĂštres d'exĂ©cution spĂ©cifiques existe. Obtenir et sauvegarder son plan d'exĂ©cution en utilisant EXPLAIN est une tĂąche Ă©lĂ©mentaire.
De plus, en utilisant l'instruction EXPLAIN (COSTS FALSE), on peut obtenir un squelette de plan, qui sera utilisé pour obtenir la valeur de hachage du plan, ce qui aidera dans l'analyse ultérieure de l'historique des modifications du plan d'exécution.
Obtenir le modÚle du plan d'exécution

  --Expliquer le plan d'exécution--
  EXECUTE 'SELECT dblink_connect(''LINK1'',''host='||database_rec.host||' dbname='||database_rec.name||' user=DATABASE password='||database_rec.owner_pwd||' '')'; 
  
  log_explain_plan = ARRAY ( SELECT * FROM dblink('LINK1', 'EXPLAIN '||log_query ) AS t (plan text) );
  log_plan_wo_costs = ARRAY ( SELECT * FROM dblink('LINK1', 'EXPLAIN ( COSTS FALSE ) '||log_query ) AS t (plan text) );
    
  PERFORM dblink_disconnect('LINK1');

Option n°3 — Utiliser le journal des requĂȘtes pour la surveillance

Étant donnĂ© que les mĂ©triques de performance ne sont pas configurĂ©es sur le texte de la requĂȘte, mais sur son ID, il est nĂ©cessaire de relier les requĂȘtes du fichier journal aux requĂȘtes pour lesquelles les mĂ©triques de performance sont configurĂ©es.
Ne serait-ce que pour avoir une heure exacte d'apparition de l'incident de performance.

Ainsi, en cas d'incident de performance pour l'ID de la requĂȘte, il y aura un lien vers la requĂȘte spĂ©cifique avec des valeurs de paramĂštres spĂ©cifiques et le temps d'exĂ©cution exact et la durĂ©e de la requĂȘte. Obtenir cette information en utilisant uniquement la vue pg_stat_statements — n'est pas possible.
Trouver le queryid de la requĂȘte et mettre Ă  jour l'enregistrement dans la table log_query

SELECT * 
	INTO log_query_rec
	FROM log_query
	WHERE query_md5hash = log_md5hash AND timepoint = log_timepoint ; 
	
	log_query_rec.query=regexp_replace(log_query_rec.query,';+','','g');
	
	FOR pg_stat_history_rec IN
	 SELECT 
      queryid ,
	  query 
	 FROM 
       pg_stat_db_queries 
     WHERE  
	   database_id = log_database_id AND
       queryid is not null 
	LOOP
	  pg_stat_query = pg_stat_history_rec.query ; 
	  pg_stat_query=regexp_replace(pg_stat_query,'n+',' ','g');
	  pg_stat_query=regexp_replace(pg_stat_query,'t+',' ','g');
	  pg_stat_query=regexp_replace(pg_stat_query,' +',' ','g');
	  pg_stat_query=regexp_replace(pg_stat_query,'$.','%','g');
	
	  log_query_text = trim(trailing ' ' from log_query_rec.query);
	  pg_stat_query_text = pg_stat_query; 
	  
	  --SELECT log_query_rec.query like pg_stat_query INTO found_flag ; 
	  IF (log_query_text LIKE pg_stat_query_text) THEN
		found_flag = TRUE ;
	  ELSE
		found_flag = FALSE ;
	  END IF;	  
	  
	  
	  IF found_flag THEN
	    
		UPDATE log_query SET queryid = pg_stat_history_rec.queryid WHERE query_md5hash = log_md5hash AND timepoint = log_timepoint ;
		activity_string = 	' updated queryid = '||pg_stat_history_rec.queryid||
		                    ' for log_query with id = '||log_query_rec.id		                    
		   				    ;						
					
	    RAISE NOTICE '%',activity_string;	
		EXIT ;
	  END IF ;
	  
	END LOOP ;

Postface

La mĂ©thode dĂ©crite a finalement trouvĂ© son application dans le systĂšme de surveillance des performances des requĂȘtes PostgreSQL en dĂ©veloppement, permettant d'avoir plus d'informations pour l'analyse lors de la rĂ©solution des incidents de performance des requĂȘtes.

Bien sûr, d'aprÚs mon point de vue personnel, il faudra encore travailler sur l'algorithme de sélection et de modification de la taille de la portion chargée. La tùche n'est pas encore résolue dans l'ensemble. Cela sera probablement intéressant.

Mais c'est une toute autre histoire 


Source : habr.com

Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS đŸ”„ Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster