Showing posts with label prise. Show all posts
Showing posts with label prise. Show all posts

Friday, March 07, 2014

How to optimize a Teradata query?

Sponsored by PRISE Ltd.
www.prisetools.com

Teradata SQL optimization techniques

Introduction

The typical goal of an SQL optimization is to get the result (data set) with less computing resources consumed and/or with shorter response time. We can follow several methodologies depending on our experience and studies, but at the end we have to get the answers for the following questions:
  • Is the task really heavy, or just the execution of the query is non-optimal?
  • What is/are the weak point(s) of the query execution?
  • What can I do to make the execution optimal?

Methodologies

The common part of the methodologies that we have to understand - more or less - what is happening during the execution. The more we understand the things behind the scenes the more we can feel the appropriate point of intervention. One can start with the trivial stuff: collect some statistics, make indices, and continue with query rewrite, or even modifying the base table structures.

What is our goal?

First of all we should branch on what do we have to do:
  1. Optimize a specific query that has been running before and we have the execution detail info
    Step details clearly show where were the big resources burnt
  2. In general, optimize the non optimal queries: find them, solve them
    Like a.,but first find those queries, and then solve them one-by-one
  3. Optimize a query, that has no detailed execution info, just the SQL (and "explain")
    Deeper knowledge of the base data and "Teradata way-of-thinking" is required, since no easy and trustworthy resource peak-detecting is available. You have to imagine what will happen, and what can be done better

Optimization in practice

This section describes the case b., and expects available detailed DBQL data.
In this post I will not attach example SQL-s, because I also switched to use PRISE Tuning Assistant for getting all the requested information for performance tuning, instead of writing complex SQL queries and making heaps of paper notes.

Prerequisites

My opinion is that DBQL (DataBase Query Logging) is the fundamental basis of a Teradata system performance management - from SQL optimization point of view. I strongly recommend to switch DBQL comprehensively ON (SQL, Step, Explain, Object are important, excluding XML, that is huge, but actually has not too much extra), and use daily archiving from the online tables - just follow Teradata recommendation.

Finding good candidate queries

DBQL is an excellent source for selecting "low hanging fruits" for performance tuning. The basic rule: we can gain big save on expensive items only, let's focus on the top resource consuming queries first. But what is high resource consumption? I usually check top queries by one or more of these properties:
  • Absolute CPU (CPU totals used by AMPs)
  • Impact CPU (CPU usage corrected by skewness)
  • Absolute I/O (I/O totals used by AMPs)
  • Impact I/O   (Disk I/O usage corrected by skewness)
  • Spool usage
  • Run duration
PRISE Tuning Assistant supplies an easy to use and quick search function for that:



Finding weak point of a query

Examining a query begins with the following steps:
  • Does it have few or many "peak steps", that consume much resources? 
    • Which one(s)?
    • What type of operations are they?
  • Does it have high skewness?
    Bad parallel efficiency, very harmful
  • Does it consume extreme huge spool?
    Compared to other queries...
PRISE Tuning Assistant again.
Check the yellow highlights in the middle, those are the top consuming steps:

    Most of the queries will have one "peak step", that consumes most of the total resources. Typical cases:
    • "Retrieve step" with redistribution
      Large number of rows and/or skewed target spool
    • "Retrieve step" with "duplication-to-all-AMPs"
      Large number of rows duplicated to all AMPs
    • Product join
      Huge number of comparisons: N * M
    • Merge or Hash join
      Skewed base or prepared (spool) data
    • OLAP function
      Large data set or skewed operation
    • Merge step
      Skewness and/or many hash collisions
    • Any kind of step
      Non small, but strongly skewed result

    What can we do?

    Teradata optimizer tries its best when produces the execution plan for a query, however it sometimes lacks proper information or its algorithms are not perfect. We - as humans - may have additional knowledge either of the data or the execution, and we can spoil the optimizer to make better decisions. Let's see our possibilities.
    • Supplement missing / refresh stale statistics
    • Drop disturbing statistics (sometimes occurs...)
    • Restructure the query
    • Break up the query, place part result into volatile table w/ good PI and put statistics on
    • Correct primary index of target / source tables
    • Build secondary/join index/indices
    • Add extra components to the query.
      You may know some additional "easy" filter that lightens the work. Eg. if you know that the join will match for only the last 3 days data of a year-covering table, you can add a date filter, which cost pennies compared to the join.
    • Restrict the result requirements to the real information demand.
      Do the end-user really need that huge amount of data, or just a record of it?

    What should we do?

    First of all, we have to find the root cause(s). Why does that specific top step consume that huge amount or resources or executes so skewed? If we find the cause and eliminate, the problem is usually solved.
    My method is the following:
    1. Find the top consuming step, and determine why it it high consumer
      • Its result is huge
      • Its result is skewed
      • Its work is huge
      • Its input(s) is/are huge
    2. Track the spool flow backwards from the top step, and find
      • Low fidelity results (row count falls far from estimated row count)
      • NO CONFIDENCE steps, specifically w/low fidelity
      • Skewed spool, specifically non small ones
      • Big duplications, specifically w/NO CONFIDENCE
    3. Find the solution
      • Supplement missing statistics, typically on PI, join fields or filter condition
        NO CONFIDENCE, low fidelity, big duplications
      • Break up the query
        Store that part result into a volatile table, where fidelity is very bad, or spool is skewed. Choose a better PI for that
      • Modify PI of the target table
        Slow MERGE step, typical hash-collision problem.
      • Eliminate product joins
      • Decompose large product join-s
      • E.T.C.
    Have a good optimization! :)

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Thursday, February 20, 2014

    DBQL analysis IV - Monitor index usage

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Analyze "Index usage" in DBQL

    Please note that the solutions found in the article works on the DBQL logs, which covers only that users' activity, for whom the logging is switched on. "Object" option in DBQL is required "on" to use the scripts provided.

    About Indices

    Teradata provides possibility of creating INDEX objects for allowing alternative access path to the data records. They are quite different structures then the good old B*Tree or Bitmap indices (common in non-MPP RDBMSes)
    The main goal of the index objects is to improve data access performance in exchange for storage and maintenance processing capacity.
    Having indices is not free (storage and maintenance resources), those ones that bring not enough gain is better being dropped.

    Index usage footprint in DBQL

    If a query uses an index for accessing the data it is declared in the "explain text", and also registered in the DBQL: dbqlobjtbl. The appearance type of indices depend on the type of index. In case of primary/secondary index only IndexId and the columns of the index are registered, while join/hash indices appear like a regular table: at database, object and column levels all.

    If the join/hash index is covering, the base table may not be listed in the DBQL objects, therefore be careful if analyze table usage from DBQLobjtbl.
    I recommend to use PRISE Tuning Assistant to easily find all type of access to a table data.

    Examples

    Primary index access (2 columns):
    Plan:
      1) First, we do a single-AMP RETRIEVE step from d01.z by way of the
         primary index "d01.z.i = 1, d01.z.j = 1" with no residual

    DBQLobjtbl:
    ...ObjectDatabaseNameObjectTableNameObjectColumnNameObjectNumObjectType...
    ...D01Zi1Idx...
    ...D01Zj1Idx...

    Secondary index access (1 column):
    Plan:
      3) We do an all-AMPs RETRIEVE step from d01.x by way of index # 4
         without accessing the base table "d01.x.j = 1" with no residual
    DBQLobjtbl:
    ...ObjectDatabaseNameObjectTableNameObjectColumnNameObjectNumObjectType...
    ...D01X
    4Idx...

    Join index:
    create join index d01.ji as sel b from d01.q primary index (b);
    select b from d01.q where a=1;
    Plan:
      1) First, we do a single-AMP RETRIEVE step from D01.JI by way of the
         primary index "D01.JI.a = 1" with no residual conditions into
    DBQLobjtbl:
    ...ObjectDatabaseNameObjectTableNameObjectColumnNameObjectNumObjectType...
    ...D01JI
    0Jix...
    ...D01JIa1Idx...
    ...D01JIb1026Col...


    Please note that
    • ObjectNum identifies the index (refers to dbc.indices.Indexnumber)
    • As many rows appeas as many columns the index has
    • Eg. in V13.10 Teradata Express the single column secondary index lacks the column name in the logs


    Analyzing DBQL data 

    Prepare data

    CREATE VOLATILE TABLE DBQLIdx_tmp1 AS (
    SELECT databasename,tablename,indexnumber,indextype,uniqueflag,Columnposition,Columnname
    , SUM (1) OVER (partition BY databasename,tablename,indexnumber ORDER BY Columnname ROWS UNBOUNDED PRECEDING) ABCOrder
    FROM dbc.indices WHERE indextype IN ('K','P','Q','S','V','H','O','I')
    ) WITH DATA
    PRIMARY INDEX (databasename,tablename,indexnumber)
    ON COMMIT PRESERVE ROWS
    ; 


    CREATE VOLATILE TABLE DBQLIdx_tmp2 AS (
    WITH RECURSIVE idxs (Databasename,Tablename,Indexnumber,Indextype,Uniqueflag,Indexcolumns,DEPTH)
    AS (
    SELECT
    databasename,tablename,indexnumber,indextype,uniqueflag,TRIM (Columnname) (VARCHAR (1000)),ABCorder
    FROM DBQLIdx_tmp1 WHERE ABCorder = 1
    UNION ALL
    SELECT
    b.databasename,b.tablename,b.indexnumber,b.indextype,b.uniqueflag,b.Indexcolumns||','||TRIM (a.Columnname),a.ABCOrder
    FROM DBQLIdx_tmp1 a
    JOIN idxs b ON a.databasename = b.databasename AND a.tablename = b.tablename AND a.indexnumber = b.indexnumber AND a.ABCOrder = b.Depth + 1
    )
    SELECT databasename db_name,tablename table_name,indextype,uniqueflag,indexcolumns
    ,indexnumber
    ,CASE WHEN uniqueflag = 'Y' AND indextype IN ('P','Q','K') THEN 'UPI'
    WHEN uniqueflag = 'N' AND indextype IN ('P','Q') THEN 'NUPI'
    WHEN uniqueflag = 'Y' AND indextype IN ('S','V','H','O') THEN 'USI'
    WHEN uniqueflag = 'N' AND indextype IN ('S','V','H','O') THEN 'NUSI'
    WHEN indextype = 'I' THEN 'O-SI'
    ELSE NULL
    END Index_code
    FROM idxs
    QUALIFY SUM (1) OVER (partition BY db_name,table_name,indexnumber ORDER BY DEPTH DESC ROWS UNBOUNDED PRECEDING) = 1
    ) WITH DATA
    PRIMARY INDEX (db_name,table_name)
    ON COMMIT PRESERVE ROWS
    ;

    UPDATE a
    FROM DBQLIdx_tmp2 a,DBQLIdx_tmp2 b
    SET Index_code = 'PK'
    WHERE a.db_name = b.db_name
    AND a.table_name = b.table_name
    AND a.Index_code = 'UPI'
    AND a.indextype = 'K'
    AND b.Index_code = 'NUPI'
    AND b.indextype <> 'K'
    ;

    Report: How many times have the indices been used?

    You may need to modify the script:
    • Date filtering (use between for interval)
    • Online/archived DBQL: use commented section for archived
    SELECT
      COALESCE(usg.objectdatabasename,idx.db_name) db
    , COALESCE(usg.objecttablename,idx.table_name) tbl
    , COALESCE(usg.ObjectNum,idx.IndexNumber) idxNo
    , idx.Index_code
    , idx.Indexcolumns Index_columnss
    , coalesce(usg.drb,0) Nbr_of_usg
    FROM
    (SELECT objectdatabasename,objecttablename,objecttype,ObjectNum,COUNT (*) drb

    --  Archived DBQL
    --  FROM dbql_arch.dbqlobjtbl_hst WHERE logdate = '2014-02-20' (date)

    --  Online DBQL
      FROM dbc.dbqlobjtbl WHERE  cast(collecttimestamp as char(10)) = '2014-02-20'
    AND objecttablename IS NOT NULL
    AND ((objecttype IN ('JIx','Hix') AND objectcolumnname IS NULL)
    OR
    (objecttype IN ('Idx'))
    )
    AND objectnum <> 1
    GROUP BY 1,2,3,4
    ) usg
    FULL OUTER JOIN
    ( SELECT db_name,table_name,Indextype,Uniqueflag,indexcolumns,Indexnumber,Index_code 

      FROM DBQLIdx_tmp2 a WHERE indextype NOT IN ('P','Q','K')
    union all
    SELECT databasename,tablename,tablekind,cast(null as char(1))

                  ,cast(null as varchar(1000)),cast(null as smallint)
    ,case when tablekind='I' then 'JIX' else 'HIX' end from dbc.tables where tablekind in ('I','N')
    ) idx ON usg.objectdatabasename = idx.db_name
    AND usg.objecttablename = idx.table_name
    AND ((usg.ObjectNum = idx.IndexNumber) or usg.objecttype IN ('JIx','Hix'))
    ORDER BY 6 DESC
    ;

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Monday, January 13, 2014

    DBQL analysis III - Monitor "collect statistics"

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Analyze "collect statistics" in DBQL

    Please note that the solutions found in the article works on the DBQL logs, which covers only that users' activity, for whom the logging is switched on. "Object" and "SQL" option in DBQL is required "on" to use the scripts provided.
    This article is applicable up to V13.10 w/o modifications, statistics handling changed from V14. 

    About Statistics

    "Statistics" is a descriptive object in the Teradata database that are used by the optimizer for transforming SQLs to effective execution plans.
    Statistics reflect the key data demographic information of one or more table column(s).
    These objects should be created and maintained, the RDBMS will not do it by itself.
    Statistics internally contain value histogram, which needs the table data (or sample) to be analyzed, which is an expensive task.

    Summarized: appropriate statistics are required for getting good and effective executon plans for SQLs, but statistics consume resources to be collected or refreshed.

    "Statistics" footprint in DBQL

    When a "statistics" is created or refreshed it is executed by an SQL command: collect statistics....
    This command will create a log entry into the DBQL if the logging is switched on.

    One can track when, which "statistics" was collected, consuming how much CPU and I/O.
    Those statements are very easy to identify in the central table:

    select * from dbc.DBQLogTbl where StatementType='Collect statistics'

    Analyzing DBQL data 

    Prepare data

    You may need to modify the script:
    • Date (interval)
    • Online/archived: use commented section
    • QueryBand: "JOB" variable is used, modify according to your ETL settings
    create volatile table DBQLStat_tmp1
    as
    (
    sel a.procId,a.QueryId,a.StartTime,a.AMPCpuTime,a.TotalIOCount

    ,case when a.querytext like '% sample %' then 'S' else 'F' end Full_Sample
    ,UserName,(FirstRespTime - StartTime) DAY(4) TO SECOND(4) AS RUNINTERVAL      
    ,(EXTRACT(DAY FROM RUNINTERVAL) * 86400 + EXTRACT(HOUR FROM RUNINTERVAL)  * 3600 + EXTRACT(MINUTE FROM RUNINTERVAL)  * 60 + EXTRACT(SECOND FROM RUNINTERVAL) ) (decimal(10,1)) Duration
    ,b.ObjectDatabaseName DatabaseName,b.ObjectTableName TableName,c.ObjectColumnName ColumnName
    ,case when d.SQLTextInfo like any ('%"PARTITION"%', '%,PARTITION %', '%,PARTITION,%', '% PARTITION,%', '% PARTITION %', '%(PARTITION,%', '%(PARTITION %', '%,PARTITION)%', '% PARTITION)%', '%(PARTITION)%') then 'Y' else 'N' end inclPartition
    ,CAST((case when index(queryband,'JOB=') >0 then  substr(queryband,index(queryband,'JOB=') ) else '' end) AS VARCHAR(500)) tmp_Q
    ,case when queryband = '' then 'N/A'
             when tmp_q = '' then '-Other'
    else CAST( (substr(tmp_Q,characters('JOB=')+1, nullifzero(index(tmp_Q,';'))-characters('JOB=')-1)) AS VARCHAR(500)) end QB_info
    ,sum(1) over (partition by a.procid,a.Queryid order by c.ObjectColumnName, a.QueryID rows unbounded preceding) Rnk
    from
    /* For achived tables
         dbql_arch.DBQLogTbl_hst       a
    join dbql_arch.DBQLObjTbl_hst      b on b.ObjectType='Tab' and a.procid=b.procid and a.QueryID=b.QueryID and a.logDate=b.logDate
    left join dbql_arch.DBQLObjTbl_hst c on c.ObjectType='Col' and a.procid=c.procid and a.QueryID=c.QueryID and a.logDate=c.logDate
    join dbql_arch.DBQLSQLTbl_hst      d on d.SQLRowNo=1       and a.procid=d.procid and a.QueryID=d.QueryID and a.logDate=d.logDate
    where a.logDate=1140113
    */
    /*end*/
    /* For online tables */
         dbc.DBQLogTbl       a
    join dbc.DBQLObjTbl      b on b.ObjectType='Tab' and a.procid=b.procid and a.QueryID=b.QueryID
    left join dbc.DBQLObjTbl c on c.ObjectType='Col' and a.procid=c.procid and a.QueryID=c.QueryID
    join dbc.DBQLSQLTbl      d on d.SQLRowNo=1       and a.procid=d.procid and a.QueryID=d.QueryID
    where cast(cast(a.starttime as char(10)) as date) = '2014-01-13' (date)
    /*end*/
    and a.StatementType='Collect statistics'
    ) with data
    primary index (procId,QueryId)
    on commit preserve rows
    ;

    create volatile table DBQLStat
    as
    (
    WITH RECURSIVE rec_tbl
    (
     procId,QueryId,StartTime,AMPCpuTime,TotalIOCount,Duration,Full_Sample,UserName,DatabaseName,TableName,QB_info,inclPartition,ColumnName,Rnk,SColumns
    )
    AS
    (
    select
     procId,QueryId,StartTime,AMPCpuTime,TotalIOCount,Duration,Full_Sample,UserName,DatabaseName,TableName,QB_info,inclPartition,ColumnName,Rnk,cast(case when ColumnName is null and inclPartition='Y' then '' else '('||ColumnName end as varchar(10000)) SColumns
    from DBQLStat_tmp1 where Rnk=1
    UNION ALL
    select
      a.procId,a.QueryId,a.StartTime,a.AMPCpuTime,a.TotalIOCount,a.Duration,a.Full_Sample,a.UserName,a.DatabaseName,a.TableName,a.QB_info,a.inclPartition,a.ColumnName,a.Rnk,b.SColumns ||','||a.ColumnName
    from DBQLStat_tmp1     a
    join rec_tbl b on a.procId=b.ProcId and a.QueryId=b.QueryID and a.Rnk=b.Rnk+1
    )
    select   procId,QueryId,StartTime,AMPCpuTime,TotalIOCount,Duration,Full_Sample,UserName,DatabaseName,TableName,QB_info,Rnk NumOfColumns
            ,case when SColumns = '' then '(PARTITION)' else SColumns || case when inclPartition='Y' then ',PARTITION)' else ')' end end StatColumns
    from rec_tbl qualify sum(1) over (partition by procid,queryid order by Rnk desc, QueryID rows unbounded preceding) = 1
    ) with data
    primary index (procid,queryid)
    on commit preserve rows
    ;

    Reports

    • How many statistics has been collected for how much resources?

    select
      UserName /*Or: DatabaseName*//*Or: Full_sample*/
    , count(*) Nbr
    , sum(AMPCpuTIme) CPU
    , sum(TotalIOCount) IO
    from DBQLStat
    group by 1
    order by 1
    ;
    • Which statistics has been collected multiple times?
      (If more days are in preapred data, frequency can be determined, erase "qualify")
    select a.*,
    sum(1) over (partition by databasename,tablename,statcolumns)  Repl
    from DBQLStat a

    /* Comment for frequency report*/
    qualify sum(1) over (partition by databasename,tablename,statcolumns) > 1

    /*end*/
    order by repl desc, databasename,tablename,statcolumns
    ;


    Sponsored by PRISE Ltd.
    www.prisetools.com

    Thursday, December 19, 2013

    Interpreting Skewness

    What does Skew metric mean?

    Overview

    You can see this word "Skewness" or "Skew factor" in a lot of places regarding Teradta: documents, applications, etc. Skewed table, skewed cpu. It is something wrong, but what does it explicitly mean? How to interpret it?

    Let's do some explanation and a bit simple maths.

    Teradata is a massive parallel system, where uniform units (AMPs) do the same tasks on that data parcel they are responsible for. In an ideal world all AMPs share the work equally, no one must work more than the average. The reality is far more cold, it is a rare situation when this equality (called "even distribution") exists.
    It is obvious that uneven distribution will cause wrong efficiency of using the parallel infrastructure.

    But how bad is the situation? Exactly that is what Skewness characterizes.

    Definitions

    Let "RESOURCE" mean the amount of resource (CPU, I/O, PERM space) consumed by an AMP.
    Let AMPno is the number of AMPs in the Teradata system.

    Skew factor := 100 - ( AVG ( "RESOURCE" ) / NULLIFZERO ( MAX ("RESOURCE") ) * 100 )

    Total[Resource] := SUM("RESOURCE")

    Impact[Resource] := MAX("RESOURCE") * AMPno

    Parallel Efficiency := Total[Resource] / Impact[Resource] * 100

    or with some transformation:

    Parallel Efficiency := 100 - Skew factor

    Analysis

    Codomain

    0 <= "Skew factor" < 100

    "Total[Resource]" <= "Impact[Resource]"

    0<"Parallel Efficiency"<=100

    Meaning

    Skew factor : This percent of the consumed real resources are wasted
    Eg. an 1Gbytes table with skew factor of 75 will allocate 4Gbytes*

    Total[Resource] :Virtual resource consumption, single sum of individual resource consumptions , measured on  AMPs as independent systems

    Impact[Resource] :Real resource consumption impacted on the parallel infrastructure

    Parallel Efficiency : As it says. Eg. Skew=80: 20%

    * Theoretically if there is/are complementary characteristics resource allocation (consumes that less resources on that AMP where my load has excess) that can compensate the parallel inefficiency from system point of view, but the probability of it tends to zero.

    Illustration



    Skew := Yellow / (Yellow + Green) * 100 [percent]


    The "Average" level indicates the mathematical average of AMP level resource consumptions (Total[Resource]), while "Peak" is maximum of AMP level resource consumptions: the real consumption from "parallel system view" (Impact[Resource])

    On finding skewed tables I will write a post later.
    PRISE Tuning Assistant helps you to find queries using CPU or I/O and helps to get rid of skewness.


    Tuesday, December 10, 2013

    DBQL analysis I. - Monitor "Top CPU consumers"

    Sponsored by PRISE Ltd.
    www.prisetools.com

    CPU usage distribution

    About DBQL

    What is it?


    DataBase Query Logging.
    It is a nice feature of Teradata RDBMS, which comprehensively logs the issued queries execution - if it is switched on.

    Configuration can be checked/administered eg. in the Teradata tools or from DBC.DBQLRuleTbl.
    Logging can be set on global/user level, and in respect of details (see DBQL tables)

    For detailed information please refer Teradata documentation of your version.

    DBQL tables


    Table Content
    DBQLogTbl Central table, 1 record for each query.
    DBQLSQLTbl Whole SQL command, broken up to 30k blocks
    DBQLStepTbl Execution steps of the query, one row for each step.
    DBQLObjTbl Objects participated in the query. Logged on different levels (db,table, column, index, etc.)
    DBQLExplainTbl English explain text, broken up to 30k blocks
    DBQLXMLTbl Explain in XML format, broken up to 30k blocks
    DBQLSummaryTbl PEs' aggregated table, which accounts on the desired level.

    DBQL tables logically organized into 1:N structure, where DBQLogTbl is the master entity and others (except DBQLSummaryTbl) are the children.
    Join fields are the ProcID and QueryId together, eg:
    ...
    from DBQLogTbl a
    join   DBQLStepTbl b on a.ProcID=b.ProcID and a.QueryID = b.QueryID
    ...
    Unfortunately PI of DBQL tables are not in sync with logical PK-FK relation in (also in latest V14.10), therefore JOIN-ed selects against online DBQL tables are not optimal.

    Cost of using DBQL

    DBQL basically consumes negligible amount of processing resources, since it has cached&batch write and generates data proportional to issued queries (flush rate is DBScontrol parameter).
    It is important to regularly purge/archive them from the DBC tables, Teradata has a recommendation for it. This ensures that PERM space consumption of the DBQL remains low.
    In an environment where ~1M SQLs are issued a day, comprehensive logging generates  ~8..10G of DBQL data daily w/o XML and Summary. Less SQLs generate proportionally less data.

    It is worth to switch on all option except XML and Summary, since the first generates huge data volume (~makes it double), and the second is similar to Acctg info. If you want to utilize them, they should be switched on, of course.

    What is it good for?

    It contains:
    • Query run time, duration
    • Consumed resources
    • Environment info (user, default db, etc)
    • SQL text
    • Explain
    • Step resource info
    • Objects involved
    • Etc.
    One can get a lot of useful aggregated and query specific tuning information, some of them I will share in the blog.

    CPU usage distribution info

    (Everything applies to I/O also, just replace CPU with I/O, AMPCPUTime with TotalIOCount...)

    Do you think Query optimization is rewarding?


    Yes, I know it is hard work to find out why is ONE query run sub-optimally, and what to do with it.

    But guess how many queries consume how many percent of the processing resources (CPU) within a whole day's workload.
    Tip it and write down for CPU%: 5%, 10%, 25% and 50%

    And now run the query below, which will result it to you. (replace the date value or maybe you have to adjust the date filtering according to local settings)

    select 'How many queries?' as "_",min(limit5) "TOP5%CPU",min(limit10) "TOP10%CPU",min(limit25) "TOP25%CPU",min(limit50) "TOP50%CPU", max(rnk) TotalQueries
    from
    (
    select
    case when CPURatio < 5.00 then null else rnk end limit5
    ,case when CPURatio < 10.00 then null else rnk end limit10
    ,case when CPURatio < 25.00 then null else rnk end limit25
    ,case when CPURatio < 50.00 then null else rnk end limit50
    ,rnk
    from
    (
    select
      sum(ampcputime) over (order by ampcputime desc ) totalCPU
    , sum(ampcputime) over (order by ampcputime desc  rows unbounded preceding) subtotalCPU
    , subtotalCPU *100.00 / totalCPU CPUratio
    , sum(1) over (order by ampcputime desc  rows unbounded preceding) rnk
    from
    (
    select *

    /* For archived DBQL
    from dbql_arch.dbqlogtbl_hst where logdate=1131201 

    and ampcputime>0
    */
    /* For online DBQL*/
    from dbc.dbqlogtbl where
    cast(cast(starttime as char(10)) as date) = '2013-12-10' (date) 

    and ampcputime>0
    ) x
    ) y
    ) z
    group by 1



    Are you surprised?
    I bet:
    • Less than 10 queries will consume 5% of the CPU
    • Less than  1% of the queries will consume 50% of the CPU
    Let's calculate.
    How much does your Teradata system cost a year? It is all for storage and processing capacity.
    If you can save eg. X% of CPU&I/O and X% storage using MVC optimization, you saved X% of the price of the Teradata system, by:
    • Improved user experience (earlier load, faster responses)
    • Resources for additional reports and applications
    • Enable postponing a very expensive Teradata hardware upgrade

    PRISE Tuning Assistant helps you to find those queries and to get the hang of how to accelerate them.

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Friday, November 29, 2013

    Accelerate skewed joins

    Sponsored by PRISE Ltd.
    www.prisetools.com

    How to "re-parallelize" skewed joins

    Case description

    Assume that we have 1M customers, 4M transactions and our top customer produce the 2.5% of all transactions.Others produce the remaining 97.5% of transactions approx. evenly.
    Scroll down to the bottom of the post for sample table and data generator SQL.

    Our task is to join a "Customer" and a "Transaction" tables on Customer_id.

    The join

    SELECT Customer_name, count(*)
    FROM Customer c
    JOIN Transact t ON c.Customer_id = t.Customer_id
    GROUP BY 1;


    We experience a pretty slow execution.
    On the ViewPoint we see that only one AMP is working, while others are not.

    What is the problem?
    There are two  separate subsets of the Transact table from "joinability" point of view:
    • "Peak" part (records of top customer(s))
      Very few customers have very much Transact records. Product join would be cost effective
    • "Even" part (records of other customers)
      Much customers have much, but specifically evenly few Transact records. Merge join would be ideal.
    Unfortunately Optimizer have to decide, only one operation type can be chosen. It will choose merge join which consumes far less CPU time.

    Execution plan looks like this:

     This query is optimized using type 2 profile T2_Linux64, profileid 21.
      1) First, we lock a distinct D_DB_TMP."pseudo table" for read on a
         RowHash to prevent global deadlock for D_DB_TMP.t.
      2) Next, we lock a distinct D_DB_TMP."pseudo table" for read on a
         RowHash to prevent global deadlock for D_DB_TMP.c.
      3) We lock D_DB_TMP.t for read, and we lock D_DB_TMP.c for read.
      4) We do an all-AMPs RETRIEVE step from D_DB_TMP.t by way of an
         all-rows scan with a condition of ("NOT (D_DB_TMP.t.Customer_ID IS
         NULL)") into Spool 4 (all_amps), which is redistributed by the
         hash code of (D_DB_TMP.t.Customer_ID) to all AMPs.  Then we do a
         SORT to order Spool 4 by row hash.  The size of Spool 4 is
         estimated with low confidence to be 125 rows (2,125 bytes).  The
         estimated time for this step is 0.01 seconds.
      5) We do an all-AMPs JOIN step from Spool 4 (Last Use) by way of a
         RowHash match scan, which is joined to D_DB_TMP.c by way of a
         RowHash match scan.  Spool 4 and D_DB_TMP.c are joined using a
         merge join, with a join condition of ("D_DB_TMP.c.Customer_ID =
         Customer_ID").  The result goes into Spool 3 (all_amps), which is
         built locally on the AMPs.  The size of Spool 3 is estimated with
         index join confidence to be 125 rows (10,375 bytes).  The
         estimated time for this step is 0.02 seconds.
      6) We do an all-AMPs SUM step to aggregate from Spool 3 (Last Use) by
         way of an all-rows scan , grouping by field1 (
         D_DB_TMP.c.Customer_name).  Aggregate Intermediate Results are
         computed globally, then placed in Spool 5.  The size of Spool 5 is
         estimated with no confidence to be 94 rows (14,758 bytes).  The
         estimated time for this step is 0.02 seconds.
      7) We do an all-AMPs RETRIEVE step from Spool 5 (Last Use) by way of
         an all-rows scan into Spool 1 (all_amps), which is built locally
         on the AMPs.  The size of Spool 1 is estimated with no confidence
         to be 94 rows (8,742 bytes).  The estimated time for this step is
         0.02 seconds.
      8) Finally, we send out an END TRANSACTION step to all AMPs involved
         in processing the request.
      -> The contents of Spool 1 are sent back to the user as the result of
         statement 1.  The total estimated time is 0.07 seconds.

    How to identify


    If you experience extremely asymmetric AMP load you can suspect on this case.
    Find highly skewed JOIN steps in the DBQL (set all logging options on):

    select top 50
    a.MaxAMPCPUTime * (hashamp()+1) / nullifzero(a.CPUTime) Skw,a.CPUTime,a.MaxAMPCPUTime * (hashamp()+1) CoveringCPUTime,
    b.*
    from dbc.dbqlsteptbl a
    join dbc.dbqlogtbl b on a.procid=b.procid and a.queryid=b.queryid
    where
    StepName='JIN'
    and CPUtime > 100
    and Skw > 2
    order by CoveringCPUTime desc;




    (Note: Covering CPU time is <No-of-AMPs> * <Max AMP's CPU time>. Virtually this amount of CPU is consumed because asymmetric load of the system)

    Or if you suspect a specific query, check the demography of the join field(s) in the "big" table:

    SELECT TOP 100 <Join_field>, count(*) Nbr
    FROM <Big_table> GROUP BY 1 ORDER BY 2 DESC;


    If the top occurences are spectacularly larger than others (or than average) the idea likely matches.


    Solution

    Break the query into two parts: join the top customer(s) separately, and then all others. Finally union the results. (Sometimes additional modification also required if the embedding operation(s) - the group by here - is/are not decomposable on the same parameter.)
    First we have to identify the top customer(s):

    SELECT TOP 5 Customer_id, count(*) Nbr
    FROM Transact GROUP BY 1 ORDER BY 2 DESC;

    Customer_id          Nbr
    ------------------------------
              345       100004
         499873                4
         677423                4
         187236                4
           23482                4
         
    Replace the original query with his one:

    SELECT Customer_name, count(*)
    FROM Customer c
    JOIN Transact t ON c.Customer_id = t.Customer_id
    where t.Customer_id in (345)  

    /*
       ID of the top Customer(s). 
       If more customers are salient, list them, but max ~5
    */
    GROUP BY 1
    UNION ALL
    SELECT Customer_name, count(*)
    FROM Customer c
    JOIN Transact t ON c.Customer_id = t.Customer_id
    where t.Customer_id not in (345)  -- Same customer(s)
    GROUP BY 1
    ;

    Be sure that Customer.Customer_id, Transact.Transact_id and Transact.Customer_id have statistics!

    Rhis query is more complex, has more steps, scans Transact table 2 times, but runs much faster, you can check it.
    But why? And how to determine which "top" customers worth to be handled separately?
    Read ahead.

    Explanation

    Calculation


    Let's do some maths:
    Assume that we are on a 125 AMP system.
    Customer table contains 1M records with unique ID.
    We have ~4.1M records in the Transact table, 100k for the top customer (ID=345), and 4 for each other customers. This matches the 2.5% we assumed above.

    If the  Transact table is redistributed on hash(Customer_id) then we will get ~33k records on each AMPs, excluding AMP(hash(345)). Here we'll get ~133k (33k + 100K).
    That means that this AMP will process ~4x more data than others, therefore runs 4x longer.
    With other words in 75% of this JOIN step's time 124 AMPs will DO NOTHING with the query.

    Moreover the preparation and subsequent steps are problematic also: the JOIN is prepared by a redistribution which produces a strongly skewed spool, and the JOIN's result stays locally on the AMPs being skewed also.

    Optimized version

    This query will consume moderately more CPU, but it is distributed evenly across the AMPs, utilizing the Teradata's full parallel capability.
    It contains a product join also, but is it no problem it joins 1 records to the selected 100k records of Transacts, that will be lightning fast.

    All

    Look at the execution plan of the broken-up query:


     This query is optimized using type 2 profile T2_Linux64, profileid 21.
      1) First, we lock a distinct D_DB_TMP."pseudo table" for read on a
         RowHash to prevent global deadlock for D_DB_TMP.t.
      2) Next, we lock a distinct D_DB_TMP."pseudo table" for read on a
         RowHash to prevent global deadlock for D_DB_TMP.c.
      3) We lock D_DB_TMP.t for read, and we lock D_DB_TMP.c for read.
      4) We do a single-AMP RETRIEVE step from D_DB_TMP.c by way of the
         unique primary index "D_DB_TMP.c.Customer_ID = 345" with no
         residual conditions into Spool 4 (all_amps), which is duplicated
         on all AMPs.  The size of Spool 4 is estimated with high
         confidence to be 125 rows (10,625 bytes).  The estimated time for
         this step is 0.01 seconds.
      5) We do an all-AMPs JOIN step from Spool 4 (Last Use) by way of an
         all-rows scan, which is joined to D_DB_TMP.t by way of an all-rows
         scan with a condition of ("D_DB_TMP.t.Customer_ID = 345").  Spool
         4 and D_DB_TMP.t are joined using a product join, with a join
         condition of ("Customer_ID = D_DB_TMP.t.Customer_ID").  The result
         goes into Spool 3 (all_amps), which is built locally on the AMPs.
         The size of Spool 3 is estimated with low confidence to be 99,670
         rows (8,272,610 bytes).  The estimated time for this step is 0.09
         seconds.
      6) We do an all-AMPs SUM step to aggregate from Spool 3 (Last Use) by
         way of an all-rows scan , grouping by field1 (
         D_DB_TMP.c.Customer_name).  Aggregate Intermediate Results are
         computed globally, then placed in Spool 5.  The size of Spool 5 is
         estimated with no confidence to be 74,753 rows (11,736,221 bytes).
         The estimated time for this step is 0.20 seconds.
      7) We execute the following steps in parallel.
           1) We do an all-AMPs RETRIEVE step from Spool 5 (Last Use) by
              way of an all-rows scan into Spool 1 (all_amps), which is
              built locally on the AMPs.  The size of Spool 1 is estimated
              with no confidence to be 74,753 rows (22,052,135 bytes).  The
              estimated time for this step is 0.02 seconds.
           2) We do an all-AMPs RETRIEVE step from D_DB_TMP.t by way of an
              all-rows scan with a condition of ("D_DB_TMP.t.Customer_ID <>
              3454") into Spool 9 (all_amps), which is redistributed by the
              hash code of (D_DB_TMP.t.Customer_ID) to all AMPs.  The size
              of Spool 9 is estimated with high confidence to be 4,294,230
              rows (73,001,910 bytes).  The estimated time for this step is
              1.80 seconds.
      8) We do an all-AMPs JOIN step from D_DB_TMP.c by way of an all-rows
         scan with a condition of ("D_DB_TMP.c.Customer_ID <> 3454"), which
         is joined to Spool 9 (Last Use) by way of an all-rows scan.
         D_DB_TMP.c and Spool 9 are joined using a single partition hash
         join, with a join condition of ("D_DB_TMP.c.Customer_ID =
         Customer_ID").  The result goes into Spool 8 (all_amps), which is
         built locally on the AMPs.  The size of Spool 8 is estimated with
         low confidence to be 4,294,230 rows (356,421,090 bytes).  The
         estimated time for this step is 0.72 seconds.
      9) We do an all-AMPs SUM step to aggregate from Spool 8 (Last Use) by
         way of an all-rows scan , grouping by field1 (
         D_DB_TMP.c.Customer_name).  Aggregate Intermediate Results are
         computed globally, then placed in Spool 10.  The size of Spool 10
         is estimated with no confidence to be 3,220,673 rows (505,645,661
         bytes).  The estimated time for this step is 8.46 seconds.
     10) We do an all-AMPs RETRIEVE step from Spool 10 (Last Use) by way of
         an all-rows scan into Spool 1 (all_amps), which is built locally
         on the AMPs.  The size of Spool 1 is estimated with no confidence
         to be 3,295,426 rows (972,150,670 bytes).  The estimated time for
         this step is 0.32 seconds.
     11) Finally, we send out an END TRANSACTION step to all AMPs involved
         in processing the request.
      -> The contents of Spool 1 are sent back to the user as the result of
         statement 1.  The total estimated time is 11.60 seconds.


    Sample structures

    The table structures (simplified for the example):


    CREATE TABLE Customer
    (
      Customer_ID   INTEGER
    , Customer_name VARCHAR(200)
    )
    UNIQUE PRIMARY INDEX (Customer_id)
    ;

    insert into Customer values (1,'Cust-1');
    Run 20x: 

    insert into Customer select mx + sum(1) over (order by Customer_id rows unbounded preceding) id, 'Cust-' || trim(id) from Customer cross join (select max(Customer_id) mx from Customer) x;

    collect statistics using sample on customer column (Customer_id);

    CREATE TABLE Transact
    (
      Transaction_ID   INTEGER
    , Customer_ID      INTEGER
    )
    UNIQUE PRIMARY INDEX (Transaction_id)
    ;

    insert into Transact values (1,1);
    Run 22x: 

    insert into Transact select mx + sum(1) over (order by Transaction_id rows unbounded preceding) id, id mod 1000000 from Transact cross join (select max(Transaction_id) mx from Transact) x;

    insert into Transact select mx + sum(1) over (order by Transaction_id rows unbounded preceding) id, 345 from Transact t cross join (select max(Transaction_id) mx from Transact) x where t.Transaction_id < 100000;

    collect statistics using sample on Transact column (Customer_id);

    collect statistics using sample on Transact column (Transaction_id) ;

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Thursday, November 28, 2013

    Optimizing Multi Value Compression

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Teradata MVC optimization
    Techniques and effects

    What is Multi Value Compression (MVC)?

    Teradata RDBMS supports a nice feature: multi-value-compression. It enables to reduce the storage space allocated by the tables in the database, while - this is incredible - processing compressed data usually requires less resources (CPU and I/O) than the uncompressed.
    The feature needs no additional licence or hardware components.

    How does MVC work?

    I give a short summary, if you are interested in the details please refer to Teradata documentation.

    MVC can be defined in CREATE TABLE DDL or later added/modified by ALTER TABLE statements. User must define a 1..255 element list of values for each compressable columns. Those will be stored as compressed value, while others will be uncompressed.
    If a column is compressed, each row has an additional area of 1..8 bits allocated (if N value is listed: upper(log2(N)) bits will be allocated). One of bit combinations means that the value is uncompressed (and allocates its corresponding space within the row layout), but all others mean compressed value, which will not allocate the value's place in the row.
    The compress bits are allocated in every rows regardless the actual value is compressed or not.
    Compress bits are "compacted", eg.: 3 + 8 + 4 = 15 compress bits will allocate 2 bytes with only 1 wasted bit instead of 3 byte aligned values.
    The value belonging to each bit combinations are stored in the table header.

    Multi Value Compression is:
    • Column level
      Have to be defined on each applicable columns of a table separately
    • "Manual"
      You have to calculate which values are worth to compress - Teradata gives no automatism
    • Static
      Once you defined the values it will not adapt to the changing conditions by itself
    It is obvious that the current optimal settings of the compression depends on the data demography and the applied data types. Optimal setting may be different later, when data demography may be different.

    Summary of most important properties of MVC once again:
    • Can be defined in the CREATE TABLE statement
    • Can be applied or modified later in an ALTER TABLE statement
    • Must be set on COLUMN level
    • Optimal value list must be calculated by you
    • Optimal settings may change in time. Optimize regularly.

    Storage effects

    Using MVC tables will allocate less PERM space, as can be calculated - simple.
    What about the increment?
    The table sizes usually grow along the time as more and more data is generated. The change in growth speed depands on the volatility of data demography. If it is stable then the growth speed will drop by the rate of compression. If typical values change in time than growth will not drop, or may speed up in extreme cases. However theese cases are when regular optimization is neccessary.

    The growth look like this in stable demography cases:




    Performance effects

    It is a key question - what have to be payed for less storage

    It is obvious that compression process requires resources during both compress and decompress phase.However there are processing gains also, which usually dominate the costs. How?

    Compressed table will reside in proportionally less data blocks, therefore data fetching requires less I/O operations. In addition moving data in-memory (during processing) requires less CPU cycles.
    While SELECTing table data usually small fragment of the row is used, and not used coulmns will not be decompressed.
    Caching is a CPU intensive operation also, which is more effective if less data blocks are processed.
    Compression helps tables to be treated as "small enough to cache 100% into memory", which results more effective execution plans.

    Summary:
    • INSERT into a compressed table usually consume more CPU by 10..50% (only final step!)
    • SELECT usually cost no more, or less CPU than at uncompressed tables
    • SELECT and INSERT usually cost proportionally less I/O like the compression ratio
    • System level CPU and I/O usage usually drops by 5..10% (!) when compressing the medium and big tables of the system (caused by more effective caching)

    How to set MVC?

    Setting up the MVC compression on a single table should consist of the following 4 steps:
    1. Analyze the data demography of each compressible columns of the table *1.
    2. Calculate the optimal compress settings for the columns. Notice that
      •   Optimum should be calculated not on separated columns, but on table level, since compress bits are packed into whole bytes.
      •   The more values are listed as compressed, the more overhead is on compress. Proper mathematical formula is to be used for calculating the optimum. *2
      •   Take care of the exceptions: PI / FK / etc.columns and some data types are not compressible (varies in different Teradata versions).
    3. Assemble the corresponding scripts
      CREATE TABLE DDL + INSERT SELECT + RENAME / ALTER TABLE DDL
    4. Implement the compression by running the script
       Concern to take good care of data protection like: backups, locking, documenting.
    *1 Simplified sample: 
         select top 256 <columnX> , count(*), avg(<length(columnX)>) from <table> group by 1 order by 2 desc; for each columns
     *2 About the algorithm: It is a maximum-seeking function (n) based on the expression of gains when specific TOP {(2^n)-1} frequent values are compressed. The expression is far more complex to discuss here because different datatypes, exceptions and internal storing constructions.

     One time or regular?

    Optimal MVC setting is valid for a specific point in time, since your data changes along your business. The daily change is usually negligible, but it accumulates.
    Practice shows that it is worth to review compress settings every 3..6 months, and continually optimize new tables, couple of weeks after coming into production.


    Estimate how much space and processing capacity is lost if compress optimization is neglected!

     

    Solution in practice

    There are "magic excels" on the net, which can calculate the optimal settings if you load the data demography, but it requires lots of manual work in addition (Running the calculations, DDL assembling, transformation script writing, testing, etc.)
     
    If you want a really simple solution, try PRISE Compress Wizard , that supplies a comprehensive solution:
    • Assists to collect good candidate tables to compress
    • Automatically analyses the tables, and gives feedback:
      • How much space can be saved by compress
      • What is the current compress ratio (if there is compress already applied)
      • How much resources were used for analysis
      • What is the optimal structure
    • Generates transforming script (+ checks, lock, logging) along with
      • Backup (arcmain)
      • Revert process (for safety and documentation)
      • Reverse engineering (for E/R documentation update)
    • Log implementation
      •  Reflect achieved space saving: success measurement
      •  Report used CPU and I/O resources for transformation

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Wednesday, September 25, 2013

    Boost slow (LEFT/RIGHT) OUTER JOINs

    Sponsored by PRISE Ltd.
    www.prisetools.com

    How to optimize slow OUTER JOINs

    Case description


    We have a (LEFT or RIGHT) OUTER JOIN, and it runs a long time while causing skewed CPU / Spool usage. In practice most of the time during the query execution only 1 AMP will work, while others have nothing to do, causing poor parallel efficiency.

    How to identify

    The query typically runs long time, contains a "MERGE JOIN" step in the Explain description, and that  step consumes skewed CPU consumption (MaxAMPCPUTime * Number-of-AMPS >> AMPCPUTime) and lasts long.

    In the DBQL you should find skewed, high CPU usage queries (
    dbc.DBQLogtbl.MaxAMPCPUTime * (hashamp()+1)  / nullifzero(dbc.DBQLogtbl.AmpCPUTime) > 1.2 and dbc.DBQLogtbl.AMPCPUTime > 1000 , depends on system size) which also has "left outer join" expression in the execution plan text (dbc.DBQLExplaintbl.ExplainText like '%left outer joined using a merge join%')
    This is only an approximation since the skewness causing step may be a different one.

    PRISE Tuning Assistant supplies easy-to-use GUI based search function.

    Explanation


    Let's assume that we outer join Table1 and Table2 on a condition that causes no product join (merge join instead), eg.:

    select Table1.x,Table2.y
    from Table1
    LEFT JOIN Table2 on Table1.x = Table2.x
    ...


    If Table2 is not a "small" table, Teradata optimizer will choose to "equi-distribute" (place matchable records on the same AMP) the two tables on the join field(s), in our case: Table1.x and Table2.x respectively.
    If Table1.x contains significant percentage of NULLs, then the distribution will be skewed, since all "x is NULL" records will get to the same AMP.
    We know that the NULL value never results in a join match, so those records are useless to examine, but they have to appear in the resultset, since it is an OUTER JOIN.


    Solution

    Let's handle the Table1 into two separate subsets: NULL(x) and NotNULL(x), and modify the select this way:

    select Table1.x,Table2.y
    from Table1
    LEFT INNER JOIN Table2 on Table1.x = Table2.x
    where Table1.x is not null -- This will eliminate skewed spool
    UNION ALL
    select Table1.x,NULL
    from Table1
    where Table1.x IS NULL;



    Practical example:
    Some of our transactions are contributed by an operator, in this case OpID is filled, else null. We would like  to query the number of transactions by operators including the non-contributed ones. Most of the transactions are non contributed ones (OpID is null).

    select
     
    a.Transaction_id, b.OperatorName as ContribOpName
    from TransactionT a
    LEFT JOIN OperatorT b on a.OpID = b.OpID


    Optimized form:

    select
      a.Transaction_id, b.OperatorName as ContribOpName
    from TransactionT a
    LEFT JOIN OperatorT b on a.OpID = b.OpID
    where a.OpID is not null
    UNION ALL
    select
      a.Transaction_id, NULL as ContribOpName
    from TransactionT
    where OpID is null;



    The execution will not cause a skewed CPU / Spool, because the those records of Table1 that caused peak ( x is NULL ) are excluded from processing of the join.
    The second part will supply the "x is NULL" records to the result set without join processing.

    The tradeoff is two full scans and a UNION ALL operation, which are comparably much less cost than a strongly skewed redistribution and a JOIN processing.

    What's next

    Next post will discuss unexpectedly slow INSERTs (hash collision).

    Sponsored by PRISE Ltd.
    www.prisetools.com

    Monday, September 23, 2013

    Accelerate PRODUCT JOIN by decomposition

    Sponsored by PRISE Ltd.
    www.prisetools.com

    How to optimize slow product joins

    Case description

    There are SQL queries that cannot be executed any other way, but using product join method, because of the content of join condition. Eg:
    • OR predicate
        Examle:
          ON (a.x = b.x OR a.y = b.y)
    • BETWEEN / LIKE operator
        Examples:
          ON (a.x LIKE b.y)
          ON (a.x LIKE b.y || '%')
          ON (a.b between b.y and b.y)
      
    • Comparison (=) of different datatype fields
        Example (a.x is INTEGER, b.x is VARCHAR)
          ON (a.x = b.x)
    • Arithmetic expression usage
        Example
          ON (a.x = b.x + 1)
    • SQL function (eg. substr() ) or UDF usage
        Example
          ON (substr(a.x,1,characters(b.x)) = b.x) 
    Product join is a technique when the execution will match each record combinations from the two joinable tables and evaluates the join condition on each of them. Product join usually causes huge CPU consumption and long response time.


    How to identify

    The query typically runs long time, contains a "PRODUCT JOIN" step in the Explain description, and that  step consumes high AMPCPUTime and lasts long. Those queries usually have >>1 LHR index (Larry Higa Ratio, showing the CPU and I/O rate), typicall 10s, 100s or more.

    In the DBQL you should find high CPU usage queries ( dbc.DBQLogtbl.AMPCPUTime > 1000 , depends on system size) which also has "product join" expression in the execution plan text (dbc.DBQLExplaintbl.ExplainText like '%product join%')

    PRISE Tuning Assistant supplies easy-to-use GUI based search function.

    Explanation of product join execution

    Let's assume that we join tables: Table1 (N records, bigger table) and Table2 (M records, smaller table) Join processing assumes that the matchable record pairs must reside on the same AMP. Since product join compares each Table1 records to each Table2 records, one of the tables' all records must reside on all AMPs, therefore PRODUCT JOIN is preceded by a "Duplicated to all AMPs" step of the smaller table.
    Each record pairs will be evaluated, if the JOIN condition satisfies, the result gets to the result spool, otherwise discarded.
    The number of required comparisons: (N x M), and the cost (approx. the required CPU time) of one comparison depends on the complexity of the join expression.

    Solution

    In most cases the JOIN condition of the product join satisfies only small fraction of all possible combinations. In practice we can identify an often situation:
    Significant subset of the bigger table's records will fit to a small subset of the smaller table's records.
    Telco example: Number analysis. Most of the Call records are directed to national number areas (>90%), but the number area describing contains dominantly international number regions (>80..95%). We can declare that national calls will never fit to international areas. In addition it is very simple to identify both a "Number" and a "Number area" if it is national or international.
    The base query looks like that:

    select Table1.Number,Table2.Area_code
    from Table1
    join Table2 ON Table1.Number BETWEEN Table2.Area_Start_number and Table2.Area_end_number;



    Let's decompose the query into two parts:

    select Table1.Number,Table2.Area_code
    from Table1
    join Table2 ON Table1.Number BETWEEN Table2.Area_Start_number and Table2.Area_end_number
    where substr(Table1.Number,1,1) = '1'   -- USA area code
    and substr(Table2.Area_Start_number,1,1) = '1' -- USA area code
    and substr(Table2.Area_end_number,1,1)   = '1' -- USA area code
    UNION ALL
    select Table1.Number,Table2.Area_code
    from Table1
    join Table2 ON Table1.Number BETWEEN Table2.Area_Start_number and Table2.Area_end_number
    where
    NOT (substr(Table1.Number,1,1) = '1')   -- non-USA area code
    and NOT
    (
        substr(Table2.Area_Start_number,1,1) = '1' -- non-USA area code
    and substr(Table2.Area_end_number,1,1)   = '1' -- non-USA area code
    );



    This way we added some low cost operations (full scan on tables to identify national/international) , and the const of UNIONing the results, but we eliminated lots of trivially not satisfing comparisions.

    The following figures show the processing cost, the red area represents the number of comparisons, therefore the cost:
    Figure1: Original case
    Figure2: Decomposed case


    Let's do some maths, with imagined combinations:
    N1: International calls
    N2: National calls
    M1: International area descriptions
    M2: National area descriptions
    90% of calls (N) are national (N2)
    90% of area descriptions (M) are international (M1).
    Originall we have to do N x M comparisons.
    The decomposed query must do
    ((0.9 x N) x (0.1 x M)) + ((0.1 x N) x (0.9 x M)) = 0.09 x N x M + 0.09 x N x M = 0.18 x N x M

    The optimized query will do only 18% of the original comparisons, with tradeoff
    of two full scans (I/O intensive) of the base tables and one UNION ALL-ing (low cost)
    of the results.

    In this case we will get ~4 times faster and CPU saving execution.


    Sometimes it can be worth to decompose to 3 or more process phases, depending on data.

    It is important, if there are further joins or transformations on the result data, they should be done on the UNION ALL-ed result, and should not be dupliated on the decomposed phases, due to code management reasons.

    Why can not do the Teradata PE the same?
    The decomposition requires the knowledge of the data, and will vary from query to query, which is currently out of scope and intelligence of an automatic optimizer.

    Summary

    Eliminate trivially invalid record pairs from the PRODUCT JOIN by breaking the query in more parts. 

    What's next

    Next post will discuss slow OUTER (LEFT/RIGHT) JOIN.

    Sponsored by PRISE Ltd.
    www.prisetools.com