From 9456cda4a7874e7f19e57f9c77c05e6521113cde Mon Sep 17 00:00:00 2001 From: bon0512 Date: Sat, 14 Mar 2026 17:12:20 +0900 Subject: [PATCH 1/4] =?UTF-8?q?[HSC-299]=20feat:=20memberllmContext=20job,?= =?UTF-8?q?=20step=20=EA=B5=AC=EC=84=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../MemberLlmContextJobConfig.java | 85 ++++++++++++++ .../tasklet/BuildMemberLlmContextTasklet.java | 51 +++++++++ .../tasklet/MemberLlmContextGateTasklet.java | 107 ++++++++++++++++++ .../VerifyMemberLlmContextTasklet.java | 104 +++++++++++++++++ 4 files changed, 347 insertions(+) create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/VerifyMemberLlmContextTasklet.java diff --git a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java new file mode 100644 index 0000000..9848c7b --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java @@ -0,0 +1,85 @@ +package site.holliverse.worker.batch.jobs.memberllmcontext; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; +import site.holliverse.worker.batch.jobs.memberllmcontext.tasklet.BuildMemberLlmContextTasklet; +import site.holliverse.worker.batch.jobs.memberllmcontext.tasklet.MemberLlmContextGateTasklet; +import site.holliverse.worker.batch.jobs.memberllmcontext.tasklet.VerifyMemberLlmContextTasklet; + +/** + * member_llm_context 적재 전용 배치 잡 구성. + * + * 흐름은 단순하게 세 단계로 나눈다. + * 1. Gate: 기준일과 전월 yyyymm 계산, 필수 테이블 확인 + * 2. Upsert: 회원 컨텍스트를 한 번에 계산해서 upsert + * 3. Verify: 적재 대상 수와 결과 수가 맞는지 최소 검증 + */ +@Configuration +@RequiredArgsConstructor +public class MemberLlmContextJobConfig { + + public static final String JOB_NAME = "memberLlmContextJob"; + + @Bean + public Job memberLlmContextJob( + JobRepository jobRepository, + Step memberLlmContextGateStep, + Step memberLlmContextUpsertStep, + Step memberLlmContextVerifyStep + ) { + return new JobBuilder(JOB_NAME, jobRepository) + .start(memberLlmContextGateStep) + .next(memberLlmContextUpsertStep) + .next(memberLlmContextVerifyStep) + .build(); + } + + /** + * 배치 파라미터와 실행 환경을 정리하는 선행 step. + */ + @Bean + public Step memberLlmContextGateStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + MemberLlmContextGateTasklet tasklet + ) { + return new StepBuilder("Step00_Gate", jobRepository) + .tasklet(tasklet, tx) + .build(); + } + + /** + * 실제 member_llm_context upsert SQL을 수행하는 핵심 step. + */ + @Bean + public Step memberLlmContextUpsertStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + BuildMemberLlmContextTasklet tasklet + ) { + return new StepBuilder("Step01_UpsertMemberLlmContext", jobRepository) + .tasklet(tasklet, tx) + .build(); + } + + /** + * 적재 결과 건수와 필수 컬럼을 검증하는 마무리 step. + */ + @Bean + public Step memberLlmContextVerifyStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + VerifyMemberLlmContextTasklet tasklet + ) { + return new StepBuilder("Step02_Verify", jobRepository) + .tasklet(tasklet, tx) + .build(); + } +} \ No newline at end of file diff --git a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java new file mode 100644 index 0000000..5060950 --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java @@ -0,0 +1,51 @@ +package site.holliverse.worker.batch.jobs.memberllmcontext.tasklet; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.jdbc.core.namedparam.MapSqlParameterSource; +import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; +import org.springframework.stereotype.Component; +import site.holliverse.worker.batch.common.support.SqlFileLoader; + +/** + * member_llm_context upsert SQL을 실제로 실행하는 tasklet. + * + * 설계상 이 step은 reader/processor/writer가 아니라, + * SQL 한 번으로 전체 적재를 끝내는 set-based 배치다. + */ +@Component +@RequiredArgsConstructor +public class BuildMemberLlmContextTasklet implements Tasklet { + + private static final String UPSERT_SQL_PATH = "sql/member-llm-context/upsert_member_llm_context.sql"; + + private final NamedParameterJdbcTemplate namedParameterJdbcTemplate; + private final SqlFileLoader sqlFileLoader; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + String snapshotDate = chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .getString("snapshotDate"); + + String yyyymm = chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .getString("yyyymm"); + + // SQL 내부에서 기준일과 전월 집계월을 함께 사용한다. + MapSqlParameterSource params = new MapSqlParameterSource() + .addValue("snapshotDate", snapshotDate) + .addValue("yyyymm", yyyymm); + + String upsertSql = sqlFileLoader.load(UPSERT_SQL_PATH); + namedParameterJdbcTemplate.update(upsertSql, params); + return RepeatStatus.FINISHED; + } +} \ No newline at end of file diff --git a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java new file mode 100644 index 0000000..6dc14a0 --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java @@ -0,0 +1,107 @@ +package site.holliverse.worker.batch.jobs.memberllmcontext.tasklet; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Component; + +import java.time.LocalDate; +import java.time.ZoneId; +import java.time.format.DateTimeFormatter; +import java.time.format.DateTimeParseException; +import java.util.ArrayList; +import java.util.List; + +/** + * member_llm_context 배치 실행 전 공통 값을 준비하는 tasklet. + * + * 역할: + * - snapshotDate 파라미터를 읽고 기본값을 보정한다. + * - usage_monthly 조회에 쓸 전월 yyyymm을 계산한다. + * - 실제 적재에 필요한 필수 테이블이 모두 있는지 확인한다. + */ +@Component +@RequiredArgsConstructor +public class MemberLlmContextGateTasklet implements Tasklet { + + private static final ZoneId KST = ZoneId.of("Asia/Seoul"); + private static final DateTimeFormatter YYYYMM = DateTimeFormatter.ofPattern("yyyyMM"); + + /** + * 현재 잡이 정상적으로 동작하려면 반드시 필요한 테이블 목록. + * churn 스냅샷이 생긴 이후에는 해당 테이블도 필수로 본다. + */ + private static final List REQUIRED_TABLES = List.of( + "member", + "subscription", + "product", + "mobile_plan", + "support_case", + "usage_monthly", + "user_event_features_7d", + "index_persona_snapshot", + "index_tscore_snapshot", + "churn_score_snapshot", + "member_llm_context" + ); + + private final JdbcTemplate jdbcTemplate; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + String snapshotDateParam = (String) chunkContext.getStepContext() + .getJobParameters() + .get("snapshotDate"); + + LocalDate snapshotDate = resolveSnapshotDate(snapshotDateParam); + String yyyymm = snapshotDate.minusMonths(1).format(YYYYMM); + + chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .putString("snapshotDate", snapshotDate.toString()); + + chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .putString("yyyymm", yyyymm); + + verifyRequiredTables(); + return RepeatStatus.FINISHED; + } + + private LocalDate resolveSnapshotDate(String snapshotDateParam) { + if (snapshotDateParam == null || snapshotDateParam.isBlank()) { + return LocalDate.now(KST); + } + + try { + return LocalDate.parse(snapshotDateParam); + } catch (DateTimeParseException e) { + throw new IllegalArgumentException("Invalid snapshotDate. Use yyyy-MM-dd format.", e); + } + } + + /** + * 필수 테이블이 하나라도 빠져 있으면 즉시 실패시키되, + * 운영에서 한 번에 원인을 볼 수 있도록 누락 목록 전체를 함께 보여준다. + */ + private void verifyRequiredTables() { + List missingTables = new ArrayList<>(); + for (String table : REQUIRED_TABLES) { + String regClass = jdbcTemplate.queryForObject("select to_regclass(?)", String.class, table); + if (regClass == null) { + missingTables.add(table); + } + } + + if (!missingTables.isEmpty()) { + throw new IllegalStateException("Missing required tables: " + String.join(", ", missingTables)); + } + } +} \ No newline at end of file diff --git a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/VerifyMemberLlmContextTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/VerifyMemberLlmContextTasklet.java new file mode 100644 index 0000000..465f320 --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/VerifyMemberLlmContextTasklet.java @@ -0,0 +1,104 @@ +package site.holliverse.worker.batch.jobs.memberllmcontext.tasklet; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.jdbc.core.namedparam.MapSqlParameterSource; +import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; +import org.springframework.stereotype.Component; +import site.holliverse.worker.batch.common.support.SqlFileLoader; + +import java.util.Map; + +/** + * member_llm_context 적재 결과를 최소 기준으로 검증하는 tasklet. + * + * 이 step의 목적은 "SQL이 실행됐다" 수준에서 멈추지 않고, + * 실제 결과가 우리가 기대한 형태로 들어갔는지 빠르게 확인하는 데 있다. + * + * 현재 검증 항목은 다음과 같다. + * 1. 적재 대상자 수와 실제 적재 row 수가 같은가 + * 2. PK(member_id)가 비어 있는 row가 없는가 + * 3. segment가 허용값(CHURN_RISK, UPSELL, NORMAL) 안에 있는가 + * 4. current_product_types가 null 없이 채워졌는가 + */ +@Component +@RequiredArgsConstructor +public class VerifyMemberLlmContextTasklet implements Tasklet { + + /** + * 검증 SQL 파일 경로. + * + * SQL에서 적재 대상 건수, 실제 적재 건수, null 건수, 잘못된 segment 건수를 + * 한 번에 가져와서 자바 쪽에서 판정한다. + */ + private static final String VERIFY_SQL_PATH = "sql/member-llm-context/verify_member_llm_context.sql"; + + private final NamedParameterJdbcTemplate namedParameterJdbcTemplate; + private final SqlFileLoader sqlFileLoader; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + // 검증 SQL을 읽어 현재 적재 상태를 한 번에 조회한다. + String verifySql = sqlFileLoader.load(VERIFY_SQL_PATH); + Map row = namedParameterJdbcTemplate.queryForMap(verifySql, new MapSqlParameterSource()); + + long eligibleCount = getLong(row, "eligible_count"); + long contextCount = getLong(row, "context_count"); + long nullPkCount = getLong(row, "null_pk_count"); + long invalidSegmentCount = getLong(row, "invalid_segment_count"); + long nullProductTypesCount = getLong(row, "null_product_types_count"); + + // 현재 배치 기준 적재 대상자 수와 실제 member_llm_context row 수가 달라지면 + // 중간 누락 또는 과적재가 발생한 것이므로 바로 실패시킨다. + if (eligibleCount != contextCount) { + throw new IllegalStateException(String.format( + "member_llm_context 건수 검증에 실패했습니다. 적재 대상 수=%d, 실제 적재 수=%d", + eligibleCount, + contextCount + )); + } + + // PK는 절대 비면 안 된다. + if (nullPkCount > 0) { + throw new IllegalStateException( + "member_llm_context PK(member_id) null 검증에 실패했습니다. nullPkCount=" + nullPkCount + ); + } + + // segment는 허용된 세 값만 들어가야 한다. + if (invalidSegmentCount > 0) { + throw new IllegalStateException( + "member_llm_context segment 값 검증에 실패했습니다. invalidSegmentCount=" + invalidSegmentCount + ); + } + + // current_product_types는 후속 추천/프롬프트 구성에 바로 쓰이므로 null이면 안 된다. + if (nullProductTypesCount > 0) { + throw new IllegalStateException( + "member_llm_context current_product_types null 검증에 실패했습니다. nullProductTypesCount=" + + nullProductTypesCount + ); + } + + return RepeatStatus.FINISHED; + } + + /** + * queryForMap 결과에서 숫자 값을 long으로 안전하게 꺼낸다. + * + * 검증 SQL이 바뀌었거나 예상과 다른 타입이 들어오면 + * 조용히 넘어가지 않고 즉시 실패시켜 원인을 빨리 드러낸다. + */ + private long getLong(Map row, String key) { + Object value = row.get(key); + if (value instanceof Number number) { + return number.longValue(); + } + throw new IllegalStateException( + "검증 SQL 결과를 숫자로 변환하지 못했습니다. key=" + key + ", value=" + value + ); + } +} \ No newline at end of file From b8c78da68f48fa687ffc3aba3b76d8537344b150 Mon Sep 17 00:00:00 2001 From: bon0512 Date: Sat, 14 Mar 2026 17:14:43 +0900 Subject: [PATCH 2/4] =?UTF-8?q?[HSC-299]=20feat:=20=EB=B0=B0=EC=B9=98=20sq?= =?UTF-8?q?l=20=EB=AC=B8=20=EC=9E=91=EC=84=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../upsert_member_llm_context.sql | 341 ++++++++++++++++++ .../verify_member_llm_context.sql | 25 ++ 2 files changed, 366 insertions(+) create mode 100644 src/main/resources/sql/member-llm-context/upsert_member_llm_context.sql create mode 100644 src/main/resources/sql/member-llm-context/verify_member_llm_context.sql diff --git a/src/main/resources/sql/member-llm-context/upsert_member_llm_context.sql b/src/main/resources/sql/member-llm-context/upsert_member_llm_context.sql new file mode 100644 index 0000000..3db16aa --- /dev/null +++ b/src/main/resources/sql/member-llm-context/upsert_member_llm_context.sql @@ -0,0 +1,341 @@ +-- member_llm_context를 한 번에 계산해서 upsert한다. +-- +-- 핵심 원칙: +-- 1) 적재 대상은 현재 활성 MOBILE_PLAN을 가진 회원만 본다. +-- 2) persona / tscore / churn은 회원별 최신 스냅샷 1건을 붙인다. +-- 3) 로그는 snapshotDate 기준 최근 14일 내 최신 1건만 붙인다. +-- 4) usage_monthly는 snapshotDate의 직전 달 yyyymm을 사용한다. +WITH base_members AS ( + -- member 전체에서 기본 회원 속성을 가져온다. + SELECT + m.member_id, + m.membership, + m.birth_date, + COALESCE(m.children_count, 0) AS children_count, + m.family_group_id, + m.family_role + FROM member m +), +active_subscriptions AS ( + -- snapshotDate 시점에 활성 상태인 구독만 남긴다. + -- current_subscriptions / current_product_types / 약정 계산의 공통 입력이다. + SELECT + s.subscription_id, + s.member_id, + s.product_id, + s.start_date, + s.end_date, + s.contract_end_date, + p.name AS product_name, + p.product_type::text AS product_type + FROM subscription s + JOIN product p ON p.product_id = s.product_id + WHERE s.status = TRUE + AND s.start_date::date <= CAST(:snapshotDate AS date) + AND (s.end_date IS NULL OR s.end_date::date >= CAST(:snapshotDate AS date)) +), +mobile_plan_subscription AS ( + -- LLM 컨텍스트 적재 대상은 활성 MOBILE_PLAN 보유 회원이다. + -- 회원당 모바일 플랜은 1개라고 가정하지만, 방어적으로 start_date ASC 1건을 고른다. + SELECT DISTINCT ON (a.member_id) + a.member_id, + a.subscription_id, + a.start_date, + mp.data_amount + FROM active_subscriptions a + JOIN mobile_plan mp ON mp.product_id = a.product_id + WHERE a.product_type = 'MOBILE_PLAN' + ORDER BY a.member_id, a.start_date ASC +), +family_counts AS ( + -- 같은 family_group_id를 가진 회원 수를 계산한다. + -- 자기 자신도 포함한다. + SELECT + m.family_group_id, + COUNT(*)::int AS family_group_num + FROM member m + WHERE m.family_group_id IS NOT NULL + GROUP BY m.family_group_id +), +latest_persona AS ( + -- 회원별 최신 persona_code 1건. + SELECT DISTINCT ON (s.member_id) + s.member_id, + s.persona_code + FROM index_persona_snapshot s + ORDER BY s.member_id, s.snapshot_date DESC +), +latest_tscore AS ( + -- 회원별 최신 tscore 1건. + -- segment의 upsell 판정에 사용한다. + SELECT DISTINCT ON (t.member_id) + t.member_id, + t.explore_tscore, + t.benefit_trend_tscore, + t.multi_device_tscore, + t.family_home_tscore + FROM index_tscore_snapshot t + ORDER BY t.member_id, t.snapshot_date DESC +), +latest_churn AS ( + -- 회원별 최신 churn 스냅샷 1건. + -- 현재는 churn_score / churn_tier 적재와 CHURN_RISK 분기에 사용한다. + SELECT DISTINCT ON (c.member_id) + c.member_id, + c.churn_score, + c.churn_tier + FROM churn_score_snapshot c + ORDER BY c.member_id, c.snapshot_date DESC +), +latest_logs_14d AS ( + -- 행동 로그는 최근 14일 이내 데이터만 허용하고 그중 최신 1건을 사용한다. + SELECT DISTINCT ON (l.member_id) + l.member_id, + l.product_type_clicks, + l.product_type_top_tags + FROM user_event_features_7d l + WHERE l.snapshot_date <= CAST(:snapshotDate AS date) + AND l.snapshot_date >= CAST(:snapshotDate AS date) - INTERVAL '14 days' + ORDER BY l.member_id, l.snapshot_date DESC +), +recent_counseling AS ( + -- 회원별 최신 상담 제목 최대 3개를 ' | '로 묶는다. + SELECT + x.member_id, + string_agg(x.title, ' | ' ORDER BY x.created_at DESC) AS recent_counseling + FROM ( + SELECT + sc.member_id, + sc.title, + sc.created_at, + row_number() OVER (PARTITION BY sc.member_id ORDER BY sc.created_at DESC) AS rn + FROM support_case sc + WHERE sc.title IS NOT NULL + AND btrim(sc.title) <> '' + ) x + WHERE x.rn <= 3 + GROUP BY x.member_id +), +current_subscriptions_json AS ( + -- 현재 활성 구독 전체를 JSON 배열로 만든다. + -- product type 5종을 모두 포함한다. + SELECT + a.member_id, + jsonb_agg( + jsonb_build_object( + 'product_id', a.product_id, + 'product_name', a.product_name, + 'product_type', a.product_type, + 'start_date', to_char(a.start_date::date, 'YYYY-MM-DD') + ) + ORDER BY a.start_date ASC + ) AS current_subscriptions + FROM active_subscriptions a + GROUP BY a.member_id +), +current_product_types_json AS ( + -- 현재 구독 중인 product_type 보유 여부를 고정 키 boolean map으로 만든다. + SELECT + a.member_id, + jsonb_build_object( + 'MOBILE_PLAN', COALESCE(bool_or(a.product_type = 'MOBILE_PLAN'), FALSE), + 'INTERNET', COALESCE(bool_or(a.product_type = 'INTERNET'), FALSE), + 'IPTV', COALESCE(bool_or(a.product_type = 'IPTV'), FALSE), + 'TAB_WATCH_PLAN', COALESCE(bool_or(a.product_type = 'TAB_WATCH_PLAN'), FALSE), + 'ADDON', COALESCE(bool_or(a.product_type = 'ADDON'), FALSE) + ) AS current_product_types + FROM active_subscriptions a + GROUP BY a.member_id +), +usage_base AS ( + -- 모바일 플랜의 데이터 제공량과 전월 사용량을 한곳에 모은다. + -- data_amount가 순수 NGB 형식이 아니면 비교 불가로 본다. + SELECT + ms.member_id, + CASE + WHEN ms.data_amount ~ '^[0-9]+(\.[0-9]+)?GB$' THEN regexp_replace(ms.data_amount, 'GB$', '', 'g')::numeric + ELSE NULL + END AS allowance_gb, + CASE + WHEN jsonb_typeof(um.usage_details) = 'object' + AND (um.usage_details ->> 'data_gb') ~ '^[0-9]+(\.[0-9]+)?$' + THEN (um.usage_details ->> 'data_gb')::numeric + ELSE NULL + END AS data_gb + FROM mobile_plan_subscription ms + LEFT JOIN usage_monthly um + ON um.subscription_id = ms.subscription_id + AND um.yyyymm = CAST(:yyyymm AS varchar(6)) +), +usage_metrics AS ( + -- usage ratio는 내부 계산용 원본 비율이다. + -- current_data_usage_ratio는 여기에서 *100 후 반올림해서 정수화한다. + SELECT + u.member_id, + CASE + WHEN u.allowance_gb IS NULL OR u.allowance_gb = 0 OR u.data_gb IS NULL THEN NULL + ELSE u.data_gb / u.allowance_gb + END AS usage_ratio + FROM usage_base u +), +contract_flags AS ( + -- ADDON을 제외한 plan 계열 구독 중 3개월 이내 약정 만료 여부를 계산한다. + SELECT + a.member_id, + bool_or( + a.product_type <> 'ADDON' + AND a.contract_end_date IS NOT NULL + AND a.contract_end_date::date >= CAST(:snapshotDate AS date) + AND a.contract_end_date::date < CAST(:snapshotDate AS date) + INTERVAL '3 months' + ) AS contract_expiry_within_3m + FROM active_subscriptions a + GROUP BY a.member_id +), +final_rows AS ( + -- 최종 적재 대상 1행을 회원별로 만든다. + -- 이 단계에서 컬럼별 최종 값과 segment를 모두 계산한다. + SELECT + bm.member_id, + bm.membership, + concat( + ((extract(year FROM age(CAST(:snapshotDate AS date), bm.birth_date))::int / 10) * 10), + chr(45824) + ) AS age_group, + ( + extract(year FROM age(CAST(:snapshotDate AS date), ms.start_date::date))::int * 12 + + extract(month FROM age(CAST(:snapshotDate AS date), ms.start_date::date))::int + ) AS join_months, + bm.children_count, + COALESCE(fc.family_group_num, 0) AS family_group_num, + bm.family_role, + lp.persona_code, + COALESCE(csj.current_subscriptions, '[]'::jsonb) AS current_subscriptions, + COALESCE( + cpt.current_product_types, + jsonb_build_object( + 'MOBILE_PLAN', FALSE, + 'INTERNET', FALSE, + 'IPTV', FALSE, + 'TAB_WATCH_PLAN', FALSE, + 'ADDON', FALSE + ) + ) AS current_product_types, + CASE + WHEN um.usage_ratio IS NULL THEN NULL + ELSE round(um.usage_ratio * 100)::int + END AS current_data_usage_ratio, + CASE + WHEN um.usage_ratio IS NULL THEN NULL + WHEN um.usage_ratio > 1.0 THEN 'OVER' + WHEN um.usage_ratio > 0.6 THEN 'FIT' + ELSE 'UNDER' + END AS data_usage_pattern, + lc.churn_score, + lc.churn_tier, + rc.recent_counseling, + ll.product_type_clicks, + COALESCE(ll.product_type_top_tags, '[]'::jsonb) AS recent_viewed_tags_top_3, + COALESCE(cf.contract_expiry_within_3m, FALSE) AS contract_expiry_within_3m, + CASE + WHEN lc.churn_tier = 'HIGH' THEN 'CHURN_RISK' + WHEN ( + (CASE + WHEN um.usage_ratio IS NULL THEN NULL + WHEN um.usage_ratio > 1.0 THEN 'OVER' + WHEN um.usage_ratio > 0.6 THEN 'FIT' + ELSE 'UNDER' + END) = 'OVER' + OR ( + SELECT COUNT(*) + FROM jsonb_array_elements_text(COALESCE(ll.product_type_top_tags, '[]'::jsonb)) AS t(tag) + WHERE t.tag IN ( + U&'OTT\D504\B9AC\BBF8\C5C4', + U&'\AC00\C871\ACB0\D569\BA54\C778', + U&'\AC00\C871\ACF5\C720', + U&'\D14C\B354\B9C1\C250\C5B4\B9C1' + ) + ) >= 2 + OR COALESCE(lt.benefit_trend_tscore, 0) >= 55 + OR COALESCE(lt.explore_tscore, 0) >= 55 + OR COALESCE(lt.multi_device_tscore, 0) >= 60 + OR COALESCE(lt.family_home_tscore, 0) >= 60 + ) THEN 'UPSELL' + ELSE 'NORMAL' + END AS segment + FROM base_members bm + JOIN mobile_plan_subscription ms ON ms.member_id = bm.member_id + LEFT JOIN family_counts fc ON fc.family_group_id = bm.family_group_id + LEFT JOIN latest_persona lp ON lp.member_id = bm.member_id + LEFT JOIN latest_tscore lt ON lt.member_id = bm.member_id + LEFT JOIN latest_churn lc ON lc.member_id = bm.member_id + LEFT JOIN latest_logs_14d ll ON ll.member_id = bm.member_id + LEFT JOIN recent_counseling rc ON rc.member_id = bm.member_id + LEFT JOIN current_subscriptions_json csj ON csj.member_id = bm.member_id + LEFT JOIN current_product_types_json cpt ON cpt.member_id = bm.member_id + LEFT JOIN usage_metrics um ON um.member_id = bm.member_id + LEFT JOIN contract_flags cf ON cf.member_id = bm.member_id +) +INSERT INTO member_llm_context ( + member_id, + membership, + age_group, + join_months, + children_count, + family_group_num, + family_role, + persona_code, + segment, + current_subscriptions, + current_product_types, + current_data_usage_ratio, + data_usage_pattern, + churn_score, + churn_tier, + recent_counseling, + product_type_clicks, + recent_viewed_tags_top_3, + contract_expiry_within_3m, + updated_at +) +SELECT + member_id, + membership, + age_group, + join_months, + children_count, + family_group_num, + family_role, + persona_code, + segment, + current_subscriptions, + current_product_types, + current_data_usage_ratio, + data_usage_pattern, + churn_score, + churn_tier, + recent_counseling, + product_type_clicks, + recent_viewed_tags_top_3, + contract_expiry_within_3m, + now() +FROM final_rows +ON CONFLICT (member_id) DO UPDATE SET + membership = EXCLUDED.membership, + age_group = EXCLUDED.age_group, + join_months = EXCLUDED.join_months, + children_count = EXCLUDED.children_count, + family_group_num = EXCLUDED.family_group_num, + family_role = EXCLUDED.family_role, + persona_code = EXCLUDED.persona_code, + segment = EXCLUDED.segment, + current_subscriptions = EXCLUDED.current_subscriptions, + current_product_types = EXCLUDED.current_product_types, + current_data_usage_ratio = EXCLUDED.current_data_usage_ratio, + data_usage_pattern = EXCLUDED.data_usage_pattern, + churn_score = EXCLUDED.churn_score, + churn_tier = EXCLUDED.churn_tier, + recent_counseling = EXCLUDED.recent_counseling, + product_type_clicks = EXCLUDED.product_type_clicks, + recent_viewed_tags_top_3 = EXCLUDED.recent_viewed_tags_top_3, + contract_expiry_within_3m = EXCLUDED.contract_expiry_within_3m, + updated_at = EXCLUDED.updated_at; \ No newline at end of file diff --git a/src/main/resources/sql/member-llm-context/verify_member_llm_context.sql b/src/main/resources/sql/member-llm-context/verify_member_llm_context.sql new file mode 100644 index 0000000..376eff3 --- /dev/null +++ b/src/main/resources/sql/member-llm-context/verify_member_llm_context.sql @@ -0,0 +1,25 @@ +-- member_llm_context 최소 검증 SQL. +-- +-- 현재 잡의 적재 대상은 "활성 MOBILE_PLAN 보유 회원"이므로 +-- member 전체가 아니라 eligible_count와 비교해야 한다. +SELECT + ( + SELECT COUNT(DISTINCT s.member_id) + FROM subscription s + JOIN product p ON p.product_id = s.product_id + WHERE s.status = TRUE + AND p.product_type::text = 'MOBILE_PLAN' + ) AS eligible_count, + (SELECT COUNT(*) FROM member_llm_context) AS context_count, + (SELECT COUNT(*) FROM member_llm_context WHERE member_id IS NULL) AS null_pk_count, + ( + SELECT COUNT(*) + FROM member_llm_context + WHERE segment NOT IN ('CHURN_RISK', 'UPSELL', 'NORMAL') + OR segment IS NULL + ) AS invalid_segment_count, + ( + SELECT COUNT(*) + FROM member_llm_context + WHERE current_product_types IS NULL + ) AS null_product_types_count; \ No newline at end of file From 88ff4c37b0712bc2845c008bb83987bd3c00fbdc Mon Sep 17 00:00:00 2001 From: bon0512 Date: Sat, 14 Mar 2026 17:38:33 +0900 Subject: [PATCH 3/4] =?UTF-8?q?[HSC-299]=20feat:=20=EC=9E=A1=20=EC=BB=A8?= =?UTF-8?q?=ED=94=BC=EA=B7=B8=20=EB=A6=AC=EC=8A=A4=EB=84=88=20=EC=88=98?= =?UTF-8?q?=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../jobs/memberllmcontext/MemberLlmContextJobConfig.java | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java index 9848c7b..ccefbcf 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/MemberLlmContextJobConfig.java @@ -9,6 +9,8 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.transaction.PlatformTransactionManager; +import site.holliverse.worker.batch.common.listener.BatchJobExecutionListener; +import site.holliverse.worker.batch.common.listener.BatchStepExecutionListener; import site.holliverse.worker.batch.jobs.memberllmcontext.tasklet.BuildMemberLlmContextTasklet; import site.holliverse.worker.batch.jobs.memberllmcontext.tasklet.MemberLlmContextGateTasklet; import site.holliverse.worker.batch.jobs.memberllmcontext.tasklet.VerifyMemberLlmContextTasklet; @@ -27,6 +29,9 @@ public class MemberLlmContextJobConfig { public static final String JOB_NAME = "memberLlmContextJob"; + private final BatchJobExecutionListener batchJobExecutionListener; + private final BatchStepExecutionListener batchStepExecutionListener; + @Bean public Job memberLlmContextJob( JobRepository jobRepository, @@ -35,6 +40,7 @@ public Job memberLlmContextJob( Step memberLlmContextVerifyStep ) { return new JobBuilder(JOB_NAME, jobRepository) + .listener(batchJobExecutionListener) .start(memberLlmContextGateStep) .next(memberLlmContextUpsertStep) .next(memberLlmContextVerifyStep) @@ -51,6 +57,7 @@ public Step memberLlmContextGateStep( MemberLlmContextGateTasklet tasklet ) { return new StepBuilder("Step00_Gate", jobRepository) + .listener(batchStepExecutionListener) .tasklet(tasklet, tx) .build(); } @@ -65,6 +72,7 @@ public Step memberLlmContextUpsertStep( BuildMemberLlmContextTasklet tasklet ) { return new StepBuilder("Step01_UpsertMemberLlmContext", jobRepository) + .listener(batchStepExecutionListener) .tasklet(tasklet, tx) .build(); } @@ -79,6 +87,7 @@ public Step memberLlmContextVerifyStep( VerifyMemberLlmContextTasklet tasklet ) { return new StepBuilder("Step02_Verify", jobRepository) + .listener(batchStepExecutionListener) .tasklet(tasklet, tx) .build(); } From c6eccca9a85713ccb226bf97b12365f8b433a848 Mon Sep 17 00:00:00 2001 From: bon0512 Date: Sat, 14 Mar 2026 17:43:51 +0900 Subject: [PATCH 4/4] =?UTF-8?q?[HSC-299]=20feat:=20=EC=A3=BC=EC=84=9D=20?= =?UTF-8?q?=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../tasklet/BuildMemberLlmContextTasklet.java | 12 ++++++++++- .../tasklet/MemberLlmContextGateTasklet.java | 21 +++++++++++++++++-- 2 files changed, 30 insertions(+), 3 deletions(-) diff --git a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java index 5060950..9c8a902 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/BuildMemberLlmContextTasklet.java @@ -39,11 +39,21 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon .getExecutionContext() .getString("yyyymm"); - // SQL 내부에서 기준일과 전월 집계월을 함께 사용한다. + /** + * SQL 내부에서 기준일과 전월 집계월을 함께 사용한다. + * snapshotDate는 나이/활성 구독/약정/로그 범위 계산에 쓰이고, + * yyyymm은 usage_monthly 전월 데이터 조회에 쓰인다. + */ MapSqlParameterSource params = new MapSqlParameterSource() .addValue("snapshotDate", snapshotDate) .addValue("yyyymm", yyyymm); + /** + * 예외 처리 이유: + * - SQL 파일을 읽지 못하면 실제 적재 로직 자체가 없다는 뜻이므로 바로 실패해야 한다. + * - SQL 실행 중 예외가 나면 조인 대상 테이블, 타입, 데이터 형식, 제약조건 중 하나가 어긋난 상황일 가능성이 크다. + * - 이 step은 핵심 적재 단계이므로 예외를 삼키지 않고 상위로 그대로 올려 배치를 실패시키는 편이 맞다. + */ String upsertSql = sqlFileLoader.load(UPSERT_SQL_PATH); namedParameterJdbcTemplate.update(upsertSql, params); return RepeatStatus.FINISHED; diff --git a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java index 6dc14a0..610b6fc 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/memberllmcontext/tasklet/MemberLlmContextGateTasklet.java @@ -75,6 +75,14 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon return RepeatStatus.FINISHED; } + /** + * snapshotDate를 yyyy-MM-dd 형식으로 파싱한다. + * 값이 없으면 한국 시간 기준 오늘 날짜를 사용한다. + * + * 예외 처리 이유: + * - 배치 기준일이 잘못되면 이후 usage, 로그 14일 범위, 약정 계산이 전부 틀어진다. + * - 그래서 형식이 맞지 않으면 조용히 보정하지 않고 즉시 실패시킨다. + */ private LocalDate resolveSnapshotDate(String snapshotDateParam) { if (snapshotDateParam == null || snapshotDateParam.isBlank()) { return LocalDate.now(KST); @@ -83,13 +91,20 @@ private LocalDate resolveSnapshotDate(String snapshotDateParam) { try { return LocalDate.parse(snapshotDateParam); } catch (DateTimeParseException e) { - throw new IllegalArgumentException("Invalid snapshotDate. Use yyyy-MM-dd format.", e); + throw new IllegalArgumentException( + "snapshotDate 파라미터 형식이 올바르지 않습니다. yyyy-MM-dd 형식을 사용하세요.", + e + ); } } /** * 필수 테이블이 하나라도 빠져 있으면 즉시 실패시키되, * 운영에서 한 번에 원인을 볼 수 있도록 누락 목록 전체를 함께 보여준다. + * + * 예외 처리 이유: + * - 이 잡은 여러 테이블을 동시에 조인하므로 필수 테이블 하나만 없어도 중간 step에서 실패한다. + * - Upsert step까지 갔다가 SQL 오류로 터지는 것보다 Gate 단계에서 빠르게 실패시키는 편이 원인 파악이 쉽다. */ private void verifyRequiredTables() { List missingTables = new ArrayList<>(); @@ -101,7 +116,9 @@ private void verifyRequiredTables() { } if (!missingTables.isEmpty()) { - throw new IllegalStateException("Missing required tables: " + String.join(", ", missingTables)); + throw new IllegalStateException( + "member_llm_context 배치 실행에 필요한 테이블이 없습니다: " + String.join(", ", missingTables) + ); } } } \ No newline at end of file