Salesforce Bulk API용 Openflow Connector 모니터링

Salesforce Bulk API용 Openflow Connector 모니터링

이 페이지에서는 Salesforce Bulk API용 Openflow Connector가 기록한 로그를 이벤트 테이블에서 조회해 복제 활동을 추적하는 방법을 설명해요. 커넥터가 어떤 객체를 몇 건 복제했는지 확인할 수 있습니다.

출처: Snowflake 문서

본문

Note

이 커넥터는 Snowflake Connector Terms에 의해 규율됩니다.

커넥터는 완료된 Salesforce Bulk API 작업에 대한 정보를 이벤트 테이블의 로그에 기록합니다. 이 로그를 조회해 Snowflake로 복제된 객체와 레코드 수를 추적할 수 있습니다.

이 페이지의 예시들은 OPENFLOW.TELEMETRY.EVENTS를 조회합니다. Openflow 배포가 텔레메트리를 다른 이벤트 테이블로 보낸다면 예시의 테이블 이름을 바꾸세요. 30분 시간 범위도 필요에 따라 조정하세요.

복제 활동 조회

기본적으로 커넥터는 Salesforce 객체 유형, 처리된 레코드 수, 벌크 작업 ID, 시스템 수정 타임스탬프를 기록합니다. Enable Merge Metrics가 false로 설정된 경우 다음 쿼리를 사용하세요:

WITH connector_logs AS (
  SELECT
    timestamp,
    resource_attributes:"openflow.dataplane.id"::VARCHAR AS deployment_id,
    resource_attributes:"k8s.namespace.name"::VARCHAR AS runtime_key,
    TRY_PARSE_JSON(value) AS parsed_log
  FROM OPENFLOW.TELEMETRY.EVENTS
  WHERE timestamp >= DATEADD('minutes', -30, CURRENT_TIMESTAMP())
    AND record_type = 'LOG'
    AND resource_attributes:"k8s.namespace.name"::VARCHAR LIKE 'runtime-%'
),
salesforce_logs AS (
  SELECT
    timestamp,
    deployment_id,
    runtime_key,
    parsed_log:formattedMessage::VARCHAR AS message
  FROM connector_logs
  WHERE parsed_log:loggerName::VARCHAR = 'org.apache.nifi.processors.standard.LogMessage'
    AND CONTAINS(parsed_log:formattedMessage::VARCHAR, 'SALESFORCE_BULK_API - ')
)
SELECT
  timestamp,
  deployment_id,
  runtime_key,
  TRIM(REGEXP_SUBSTR(message, 'ObjectType = ([^;]+)', 1, 1, 'e', 1)) AS object_type,
  TRY_TO_NUMBER(TRIM(REGEXP_SUBSTR(message, 'Records = ([^;]+)', 1, 1, 'e', 1))) AS records,
  TRIM(REGEXP_SUBSTR(message, 'BulkJobID = ([^;]*)', 1, 1, 'e', 1)) AS bulk_job_id,
  TRIM(REGEXP_SUBSTR(message, 'SystemModstamp = ([^;]+)', 1, 1, 'e', 1)) AS system_modstamp
FROM salesforce_logs
ORDER BY timestamp DESC;

병합 메트릭 조회

로그에 상세 레코드 수를 포함하려면 Enable Merge Metrics를 true로 설정하세요. 커넥터는 각 증분 병합 전에 추가 쿼리를 실행해 그 수치를 계산합니다. 이 쿼리는 Snowflake Warehouse에 구성된 웨어하우스를 사용합니다.

병합 메트릭은 IsDeleted 필드를 포함하는 Salesforce 객체에 대해서만 사용할 수 있습니다. 복제 활동과 병합 메트릭을 조회하려면 다음 쿼리를 사용하세요:

WITH connector_logs AS (
  SELECT
    timestamp,
    resource_attributes:"openflow.dataplane.id"::VARCHAR AS deployment_id,
    resource_attributes:"k8s.namespace.name"::VARCHAR AS runtime_key,
    TRY_PARSE_JSON(value) AS parsed_log
  FROM OPENFLOW.TELEMETRY.EVENTS
  WHERE timestamp >= DATEADD('minutes', -30, CURRENT_TIMESTAMP())
    AND record_type = 'LOG'
    AND resource_attributes:"k8s.namespace.name"::VARCHAR LIKE 'runtime-%'
),
salesforce_logs AS (
  SELECT
    timestamp,
    deployment_id,
    runtime_key,
    parsed_log:formattedMessage::VARCHAR AS message
  FROM connector_logs
  WHERE parsed_log:loggerName::VARCHAR = 'org.apache.nifi.processors.standard.LogMessage'
    AND CONTAINS(parsed_log:formattedMessage::VARCHAR, 'SALESFORCE_BULK_API - ')
)
SELECT
  timestamp,
  deployment_id,
  runtime_key,
  TRIM(REGEXP_SUBSTR(message, 'ObjectType = ([^;]+)', 1, 1, 'e', 1)) AS object_type,
  TRY_TO_NUMBER(TRIM(REGEXP_SUBSTR(message, 'Records = ([^;]+)', 1, 1, 'e', 1))) AS records,
  TRIM(REGEXP_SUBSTR(message, 'BulkJobID = ([^;]*)', 1, 1, 'e', 1)) AS bulk_job_id,
  TRIM(REGEXP_SUBSTR(message, 'SystemModstamp = ([^;]+)', 1, 1, 'e', 1)) AS system_modstamp,
  TRY_TO_NUMBER(TRIM(REGEXP_SUBSTR(message, 'ROWS_ADDED = ([^;]+)', 1, 1, 'e', 1))) AS rows_added,
  TRY_TO_NUMBER(TRIM(REGEXP_SUBSTR(message, 'ROWS_ADDED_DELETED = ([^;]+)', 1, 1, 'e', 1))) AS rows_added_deleted,
  TRY_TO_NUMBER(TRIM(REGEXP_SUBSTR(message, 'ROWS_UPDATED = ([^;]+)', 1, 1, 'e', 1))) AS rows_updated,
  TRY_TO_NUMBER(TRIM(REGEXP_SUBSTR(message, 'ROWS_DELETED = ([^;]+)', 1, 1, 'e', 1))) AS rows_deleted,
  TRY_TO_NUMBER(TRIM(REGEXP_SUBSTR(message, 'ROWS_RESTORED = ([^;]+)', 1, 1, 'e', 1))) AS rows_restored
FROM salesforce_logs
ORDER BY timestamp DESC;

메트릭의 의미는 다음과 같습니다:

Metric Initial load Incremental load
ROWS_ADDED 삭제로 표시되지 않은 로드된 레코드. 대상 테이블에 존재하지 않는 활성 소스 레코드.
ROWS_ADDED_DELETED 이미 삭제로 표시된 로드된 레코드. 대상 테이블에 존재하지 않는, 삭제로 표시된 소스 레코드.
ROWS_UPDATED 0 대상 테이블의 활성 레코드와 일치하는 활성 소스 레코드.
ROWS_DELETED 0 대상 테이블의 활성 레코드와 일치하는, 삭제로 표시된 소스 레코드.
ROWS_RESTORED 0 대상 테이블에서 삭제로 표시된 레코드와 일치하는 활성 소스 레코드.

ROWS_UPDATED는 대상의 활성 레코드와 일치하는 활성 소스 레코드를 셉니다. 개별 필드 값을 비교하지 않으며, 필드 값이 변경되었는지는 나타내지 않습니다.

로그에 병합 메트릭이 없으면 쿼리의 병합 메트릭 열이 NULL을 반환합니다. 이는 Enable Merge Metrics가 false로 설정되었거나, Salesforce 객체에 IsDeleted 필드가 없거나, 로그가 병합 메트릭을 지원하지 않는 이전 커넥터 버전에서 생성되었을 때 발생합니다.

더 알아보기 (Learn more)