From 351a9db097862bead6f912ac6b3460709ee86d71 Mon Sep 17 00:00:00 2001 From: bon0512 Date: Thu, 12 Mar 2026 01:46:47 +0900 Subject: [PATCH 1/5] =?UTF-8?q?[HSC-200]=20feat:=20=EC=9E=A1=20=EC=BB=A8?= =?UTF-8?q?=ED=94=BC=EA=B7=B8=20=EC=84=A4=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../jobs/index/WeeklyIndexJobConfig.java | 4 ++++ .../worker/batch/jobs/index/tasklet/.gitkeep | 0 .../index/tasklet/GateWeeklyIndexTasklet.java | 23 +++++++++++++++++++ 3 files changed, 27 insertions(+) create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/.gitkeep create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java b/src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java new file mode 100644 index 0000000..db4cfcd --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java @@ -0,0 +1,4 @@ +package site.holliverse.worker.batch.jobs.index; + +public class WeeklyIndexJobConfig { +} diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/.gitkeep b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java new file mode 100644 index 0000000..82d2715 --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java @@ -0,0 +1,23 @@ +package site.holliverse.worker.batch.jobs.index.tasklet; + + +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 java.time.ZoneId; +import java.time.format.DateTimeFormatter; + +public class GateWeeklyIndexTasklet implements Tasklet { + + private static final ZoneId KST = ZoneId.of("Asia/Seoul"); + private static final DateTimeFormatter YYYYMM = DateTimeFormatter.ofPattern("yyyyMM"); + + + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { + return null; + } +} From 47ad95d7509b55e1ae79ba629e818e2aff33e922 Mon Sep 17 00:00:00 2001 From: bon0512 Date: Thu, 12 Mar 2026 01:48:53 +0900 Subject: [PATCH 2/5] =?UTF-8?q?[HSC-200]=20feat:=20sql=20fileLoader=20?= =?UTF-8?q?=EC=9E=91=EC=84=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../batch/common/support/SqlFileLoader.java | 22 +++++++++++++++++++ 1 file changed, 22 insertions(+) create mode 100644 src/main/java/site/holliverse/worker/batch/common/support/SqlFileLoader.java diff --git a/src/main/java/site/holliverse/worker/batch/common/support/SqlFileLoader.java b/src/main/java/site/holliverse/worker/batch/common/support/SqlFileLoader.java new file mode 100644 index 0000000..8a87758 --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/common/support/SqlFileLoader.java @@ -0,0 +1,22 @@ +package site.holliverse.worker.batch.common.support; + +import org.springframework.core.io.ClassPathResource; +import org.springframework.stereotype.Component; +import org.springframework.util.StreamUtils; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; + +@Component +public class SqlFileLoader { + + // classpath SQL 파일을 UTF-8 문자열로 읽어 Tasklet에서 재사용한다. + public String load(String classpathLocation) { + ClassPathResource resource = new ClassPathResource(classpathLocation); + try { + return StreamUtils.copyToString(resource.getInputStream(), StandardCharsets.UTF_8); + } catch (IOException e) { + throw new IllegalStateException("SQL 파일을 읽는 중 오류가 발생했습니다: " + classpathLocation, e); + } + } +} From 15a70a92095a8feb7f0be8fcbe6c406d6c33eac8 Mon Sep 17 00:00:00 2001 From: bon0512 Date: Thu, 12 Mar 2026 01:49:34 +0900 Subject: [PATCH 3/5] =?UTF-8?q?[HSC-200]=20feat:=20tasklet=20=EC=9E=91?= =?UTF-8?q?=EC=84=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../BuildPersonaWeeklyIndexTasklet.java | 63 +++++++++++ .../tasklet/BuildRawWeeklyIndexTasklet.java | 68 ++++++++++++ .../BuildTScoreWeeklyIndexTasklet.java | 58 ++++++++++ .../tasklet/VerifyWeeklyIndexTasklet.java | 101 ++++++++++++++++++ 4 files changed, 290 insertions(+) create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java create mode 100644 src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java new file mode 100644 index 0000000..edee6da --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java @@ -0,0 +1,63 @@ +package site.holliverse.worker.batch.jobs.index.tasklet; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +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; + +/** + * 페르소나 스냅샷을 생성하는 Step. + * + * 처리 순서: + * 1) 동일 snapshotDate의 기존 persona 결과를 삭제 + * 2) index_tscore_snapshot에서 회원별 최고 T-score 지수 1개 선택 + * 3) 지수 코드 -> persona_code 매핑 후 index_persona_snapshot에 저장 + * + * 동점 규칙: + * - SQL 내부 정렬 기준(Order By)에 따라 해소 + * - 현재는 index_code 오름차순(사전순) 우선 + */ +@Component +@RequiredArgsConstructor +@Slf4j +public class BuildPersonaWeeklyIndexTasklet implements Tasklet { + + private static final String DELETE_SQL_PATH = "sql/index/delete_persona_snapshot.sql"; + private static final String INSERT_SQL_PATH = "sql/index/insert_persona_snapshot.sql"; + + private final NamedParameterJdbcTemplate namedParameterJdbcTemplate; + private final SqlFileLoader sqlFileLoader; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + // 1) Gate Step에서 확정한 snapshotDate를 컨텍스트에서 읽는다. + String snapshotDate = chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .getString("snapshotDate"); + + // 2) SQL 바인딩 파라미터를 준비한다. + MapSqlParameterSource params = new MapSqlParameterSource() + .addValue("snapshotDate", snapshotDate); + + // 3) 실행할 SQL 파일(삭제/적재)을 로드한다. + String deleteSql = sqlFileLoader.load(DELETE_SQL_PATH); + String insertSql = sqlFileLoader.load(INSERT_SQL_PATH); + + // 4) 멱등 실행을 위해 기존 데이터 삭제 후 재적재한다. + int deleted = namedParameterJdbcTemplate.update(deleteSql, params); + int inserted = namedParameterJdbcTemplate.update(insertSql, params); + + log.info("BuildPersona 완료. snapshotDate={}, deleted={}, inserted={}", + snapshotDate, deleted, inserted); + + return RepeatStatus.FINISHED; + } +} diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java new file mode 100644 index 0000000..35c040e --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java @@ -0,0 +1,68 @@ +package site.holliverse.worker.batch.jobs.index.tasklet; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +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; + +/** + * raw 지수 스냅샷을 생성하는 Step. + * + * 처리 방식: + * - 동일 snapshotDate 기존 결과 삭제 + * - 같은 기준으로 raw 결과 전체 재적재 + * + * 목적: + * - 재실행 시에도 결과를 안정적으로 덮어써 멱등성을 보장한다. + */ +@Component +@RequiredArgsConstructor +@Slf4j +public class BuildRawWeeklyIndexTasklet implements Tasklet { + + private static final String DELETE_SQL_PATH = "sql/index/delete_raw_snapshot.sql"; + private static final String INSERT_SQL_PATH = "sql/index/insert_raw_snapshot.sql"; + + private final NamedParameterJdbcTemplate namedParameterJdbcTemplate; + private final SqlFileLoader sqlFileLoader; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + // 1) Gate Step에서 저장한 기준 값을 읽는다. + String snapshotDate = chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .getString("snapshotDate"); + + String yyyymm = chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .getString("yyyymm"); + + // 2) SQL 파라미터를 구성한다. + MapSqlParameterSource params = new MapSqlParameterSource() + .addValue("snapshotDate", snapshotDate) + .addValue("yyyymm", yyyymm); + + // 3) 클래스패스 SQL 파일을 읽는다. + String deleteSql = sqlFileLoader.load(DELETE_SQL_PATH); + String insertSql = sqlFileLoader.load(INSERT_SQL_PATH); + + // 4) 기존 데이터 삭제 후 재계산 결과를 적재한다. + int deleted = namedParameterJdbcTemplate.update(deleteSql, params); + int inserted = namedParameterJdbcTemplate.update(insertSql, params); + + log.info("BuildRaw 완료. snapshotDate={}, yyyymm={}, deleted={}, inserted={}", + snapshotDate, yyyymm, deleted, inserted); + + return RepeatStatus.FINISHED; + } +} diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java new file mode 100644 index 0000000..fc5e87d --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java @@ -0,0 +1,58 @@ +package site.holliverse.worker.batch.jobs.index.tasklet; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +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; + +/** + * T-score 스냅샷을 생성하는 Step. + * + * 처리 방식: + * - 동일 snapshotDate 기존 tscore 결과 삭제 + * - raw 분포(avg/stddev) 기반으로 tscore 재계산 후 적재 + */ +@Component +@RequiredArgsConstructor +@Slf4j +public class BuildTScoreWeeklyIndexTasklet implements Tasklet { + + private static final String DELETE_SQL_PATH = "sql/index/delete_tscore_snapshot.sql"; + private static final String INSERT_SQL_PATH = "sql/index/insert_tscore_snapshot.sql"; + + private final NamedParameterJdbcTemplate namedParameterJdbcTemplate; + private final SqlFileLoader sqlFileLoader; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + // 1) 현재 스냅샷 날짜를 컨텍스트에서 읽는다. + String snapshotDate = chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .getString("snapshotDate"); + + // 2) SQL 파라미터를 구성한다. + MapSqlParameterSource params = new MapSqlParameterSource() + .addValue("snapshotDate", snapshotDate); + + // 3) SQL 파일을 로드한다. + String deleteSql = sqlFileLoader.load(DELETE_SQL_PATH); + String insertSql = sqlFileLoader.load(INSERT_SQL_PATH); + + // 4) 기존 결과를 초기화하고 새 결과를 적재한다. + int deleted = namedParameterJdbcTemplate.update(deleteSql, params); + int inserted = namedParameterJdbcTemplate.update(insertSql, params); + + log.info("BuildTScore 완료. snapshotDate={}, deleted={}, inserted={}", + snapshotDate, deleted, inserted); + + return RepeatStatus.FINISHED; + } +} diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java new file mode 100644 index 0000000..dcb7c35 --- /dev/null +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java @@ -0,0 +1,101 @@ +package site.holliverse.worker.batch.jobs.index.tasklet; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +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; + +/** + * 적재 결과를 최종 검증하는 Step. + * + * 검증 항목: + * - member 수 == raw 수 == tscore 수 == persona 수 + * - tscore 필수 컬럼 null 여부 + * - persona 필수 컬럼 null 여부 + */ +@Component +@RequiredArgsConstructor +@Slf4j +public class VerifyWeeklyIndexTasklet implements Tasklet { + + private static final String VERIFY_COUNTS_SQL_PATH = "sql/index/verify_counts.sql"; + private static final String VERIFY_TSCORE_NULLS_SQL_PATH = "sql/index/verify_tscore_nulls.sql"; + private static final String VERIFY_PERSONA_NULLS_SQL_PATH = "sql/index/verify_persona_nulls.sql"; + + private final NamedParameterJdbcTemplate namedParameterJdbcTemplate; + private final SqlFileLoader sqlFileLoader; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + // 1) 검증 대상 snapshotDate를 읽는다. + String snapshotDate = chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .getString("snapshotDate"); + + // 2) SQL 파라미터를 구성한다. + MapSqlParameterSource params = new MapSqlParameterSource() + .addValue("snapshotDate", snapshotDate); + + // 3) 검증 SQL 파일을 로드한다. + String verifyCountsSql = sqlFileLoader.load(VERIFY_COUNTS_SQL_PATH); + String verifyTScoreNullsSql = sqlFileLoader.load(VERIFY_TSCORE_NULLS_SQL_PATH); + String verifyPersonaNullsSql = sqlFileLoader.load(VERIFY_PERSONA_NULLS_SQL_PATH); + + // 4) 건수 검증: 모집단 대비 누락/중복 여부를 확인한다. + Map counts = namedParameterJdbcTemplate.queryForMap(verifyCountsSql, params); + long memberCount = getLong(counts, "member_count"); + long rawCount = getLong(counts, "raw_count"); + long tscoreCount = getLong(counts, "tscore_count"); + long personaCount = getLong(counts, "persona_count"); + + if (memberCount != rawCount || rawCount != tscoreCount || tscoreCount != personaCount) { + throw new IllegalStateException(String.format( + "건수 검증에 실패했습니다. member=%d, raw=%d, tscore=%d, persona=%d", + memberCount, rawCount, tscoreCount, personaCount + )); + } + + // 5) tscore null 검증 + Long tscoreNullCount = namedParameterJdbcTemplate.queryForObject(verifyTScoreNullsSql, params, Long.class); + if (tscoreNullCount != null && tscoreNullCount > 0) { + throw new IllegalStateException("T-score null 검증에 실패했습니다. nullCount=" + tscoreNullCount); + } + + // 6) persona null 검증 + Long personaNullCount = namedParameterJdbcTemplate.queryForObject(verifyPersonaNullsSql, params, Long.class); + if (personaNullCount != null && personaNullCount > 0) { + throw new IllegalStateException("Persona null 검증에 실패했습니다. nullCount=" + personaNullCount); + } + + log.info( + "Verify 완료. snapshotDate={}, memberCount={}, rawCount={}, tscoreCount={}, personaCount={}, tscoreNullCount={}, personaNullCount={}", + snapshotDate, + memberCount, + rawCount, + tscoreCount, + personaCount, + tscoreNullCount, + personaNullCount + ); + + return RepeatStatus.FINISHED; + } + + private long getLong(Map row, String key) { + Object value = row.get(key); + if (value instanceof Number number) { + return number.longValue(); + } + throw new IllegalStateException("숫자 값 변환에 실패했습니다. key=" + key + ", value=" + value); + } +} From 504b1506edb0ccc94b9d6b5b9ddd8f92fb63e8ab Mon Sep 17 00:00:00 2001 From: bon0512 Date: Thu, 12 Mar 2026 01:50:37 +0900 Subject: [PATCH 4/5] =?UTF-8?q?[HSC-200]=20feat:=20sql=20=EC=9E=91?= =?UTF-8?q?=EC=84=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 3 +- .../jobs/index/WeeklyIndexJobConfig.java | 93 +++++++ .../index/tasklet/GateWeeklyIndexTasklet.java | 100 ++++++- .../index/create_index_snapshot_tables.sql | 82 ++++++ .../sql/index/delete_persona_snapshot.sql | 3 + .../sql/index/delete_raw_snapshot.sql | 3 + .../sql/index/delete_tscore_snapshot.sql | 3 + .../sql/index/insert_persona_snapshot.sql | 79 ++++++ .../sql/index/insert_raw_snapshot.sql | 248 ++++++++++++++++++ .../sql/index/insert_tscore_snapshot.sql | 74 ++++++ .../resources/sql/index/verify_counts.sql | 6 + .../sql/index/verify_persona_nulls.sql | 9 + .../sql/index/verify_tscore_nulls.sql | 12 + 13 files changed, 711 insertions(+), 4 deletions(-) create mode 100644 src/main/resources/sql/index/create_index_snapshot_tables.sql create mode 100644 src/main/resources/sql/index/delete_persona_snapshot.sql create mode 100644 src/main/resources/sql/index/delete_raw_snapshot.sql create mode 100644 src/main/resources/sql/index/delete_tscore_snapshot.sql create mode 100644 src/main/resources/sql/index/insert_persona_snapshot.sql create mode 100644 src/main/resources/sql/index/insert_raw_snapshot.sql create mode 100644 src/main/resources/sql/index/insert_tscore_snapshot.sql create mode 100644 src/main/resources/sql/index/verify_counts.sql create mode 100644 src/main/resources/sql/index/verify_persona_nulls.sql create mode 100644 src/main/resources/sql/index/verify_tscore_nulls.sql diff --git a/.gitignore b/.gitignore index 530fe04..cd04363 100644 --- a/.gitignore +++ b/.gitignore @@ -12,4 +12,5 @@ pnpm-debug.log* npm-debug.log* yarn-debug.log* yarn-error.log* -.gradle-home/ \ No newline at end of file +.gradle-home/ +docs/private/ diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java b/src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java index db4cfcd..8ff86f0 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/WeeklyIndexJobConfig.java @@ -1,4 +1,97 @@ package site.holliverse.worker.batch.jobs.index; +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.index.tasklet.BuildPersonaWeeklyIndexTasklet; +import site.holliverse.worker.batch.jobs.index.tasklet.BuildRawWeeklyIndexTasklet; +import site.holliverse.worker.batch.jobs.index.tasklet.BuildTScoreWeeklyIndexTasklet; +import site.holliverse.worker.batch.jobs.index.tasklet.GateWeeklyIndexTasklet; +import site.holliverse.worker.batch.jobs.index.tasklet.VerifyWeeklyIndexTasklet; + +@Configuration +@RequiredArgsConstructor public class WeeklyIndexJobConfig { + + public static final String JOB_NAME = "weeklyIndexJob"; + + @Bean + public Job weeklyIndexJob( + JobRepository jobRepository, + Step gateStep, + Step buildRawStep, + Step buildTScoreStep, + Step buildPersonaStep, + Step verifyStep + ) { + // 주간 지수 배치 플로우: 준비 -> raw -> tscore -> persona -> 검증 + return new JobBuilder(JOB_NAME, jobRepository) + .start(gateStep) + .next(buildRawStep) + .next(buildTScoreStep) + .next(buildPersonaStep) + .next(verifyStep) + .build(); + } + + @Bean + public Step gateStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + GateWeeklyIndexTasklet tasklet + ) { + return new StepBuilder("Step00_Gate", jobRepository) + .tasklet(tasklet, tx) + .build(); + } + + @Bean + public Step buildRawStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + BuildRawWeeklyIndexTasklet tasklet + ) { + return new StepBuilder("Step01_BuildRaw", jobRepository) + .tasklet(tasklet, tx) + .build(); + } + + @Bean + public Step buildTScoreStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + BuildTScoreWeeklyIndexTasklet tasklet + ) { + return new StepBuilder("Step02_BuildTScore", jobRepository) + .tasklet(tasklet, tx) + .build(); + } + + @Bean + public Step buildPersonaStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + BuildPersonaWeeklyIndexTasklet tasklet + ) { + return new StepBuilder("Step03_BuildPersona", jobRepository) + .tasklet(tasklet, tx) + .build(); + } + + @Bean + public Step verifyStep( + JobRepository jobRepository, + PlatformTransactionManager tx, + VerifyWeeklyIndexTasklet tasklet + ) { + return new StepBuilder("Step04_Verify", jobRepository) + .tasklet(tasklet, tx) + .build(); + } } diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java index 82d2715..c6c256f 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java @@ -1,23 +1,117 @@ package site.holliverse.worker.batch.jobs.index.tasklet; - +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; 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.List; +/** + * 주간 지수 배치의 게이트 Step. + * + * 역할: + * - snapshotDate 파라미터를 읽고 유효성을 검증한다. + * - 월 기준 키(yyyymm)를 계산한다. + * - 후속 Step이 공통으로 쓰는 값을 ExecutionContext에 저장한다. + * - 필수 소스/타깃 테이블 존재 여부를 사전 검증한다. + */ +@Component +@RequiredArgsConstructor +@Slf4j public class GateWeeklyIndexTasklet implements Tasklet { private static final ZoneId KST = ZoneId.of("Asia/Seoul"); private static final DateTimeFormatter YYYYMM = DateTimeFormatter.ofPattern("yyyyMM"); + // 배치 시작 전에 반드시 존재해야 하는 테이블 목록 + private static final List REQUIRED_TABLES = List.of( + "member", + "subscription", + "product", + "addon_service", + "internet", + "iptv", + "mobile_plan", + "tab_watch_plan", + "member_coupon", + "usage_monthly", + "support_case", + "billing", + "user_event_features_7d", + "index_raw_snapshot", + "index_tscore_snapshot", + "index_persona_snapshot" + ); + private final JdbcTemplate jdbcTemplate; @Override - public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { - return null; + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + // 1) snapshotDate 파라미터를 읽는다. + String snapshotDateParam = (String) chunkContext.getStepContext() + .getJobParameters() + .get("snapshotDate"); + + // 2) 기준 날짜와 월 키를 계산한다. + LocalDate snapshotDate = resolveSnapshotDate(snapshotDateParam); + // 월 기준 지표(usage/billing/support)는 전월 확정치를 사용한다. + String yyyymm = snapshotDate.minusMonths(1).format(YYYYMM); + + // 3) 후속 Step에서 재사용하도록 컨텍스트에 저장한다. + chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .putString("snapshotDate", snapshotDate.toString()); + + chunkContext.getStepContext() + .getStepExecution() + .getJobExecution() + .getExecutionContext() + .putString("yyyymm", yyyymm); + + // 4) 필수 테이블 존재 여부를 확인한다. + verifyRequiredTables(); + + log.info("Gate 통과. snapshotDate={}, yyyymm={}", snapshotDate, yyyymm); + return RepeatStatus.FINISHED; + } + + private LocalDate resolveSnapshotDate(String snapshotDateParam) { + // 파라미터가 없으면 KST 오늘 날짜를 기본값으로 사용한다. + if (snapshotDateParam == null || snapshotDateParam.isBlank()) { + return LocalDate.now(KST); + } + + // yyyy-MM-dd 형식만 허용한다. + try { + return LocalDate.parse(snapshotDateParam); + } catch (DateTimeParseException e) { + throw new IllegalArgumentException("snapshotDate 형식이 올바르지 않습니다. yyyy-MM-dd 형식을 사용하세요.", e); + } + } + + private void verifyRequiredTables() { + // to_regclass는 테이블이 없으면 null을 반환한다. + for (String table : REQUIRED_TABLES) { + String regClass = jdbcTemplate.queryForObject( + "select to_regclass(?)", + String.class, + table + ); + + if (regClass == null) { + throw new IllegalStateException("필수 테이블이 존재하지 않습니다: " + table); + } + } } } diff --git a/src/main/resources/sql/index/create_index_snapshot_tables.sql b/src/main/resources/sql/index/create_index_snapshot_tables.sql new file mode 100644 index 0000000..407de6a --- /dev/null +++ b/src/main/resources/sql/index/create_index_snapshot_tables.sql @@ -0,0 +1,82 @@ +-- 목적: 주간 지수 배치 결과(raw/tscore/persona) 저장 테이블 생성 +-- 특이사항: +-- - snapshot_date + member_id 를 기본 키로 사용해 동일 날짜 중복 적재를 방지한다. +-- - persona 결과는 persona_type_id(FK) 기준으로 저장한다. + +CREATE TABLE IF NOT EXISTS index_raw_snapshot ( + snapshot_date date NOT NULL, + member_id bigint NOT NULL, + explore_raw numeric NOT NULL, + benefit_trend_raw numeric NOT NULL, + multi_device_raw numeric NOT NULL, + family_home_raw numeric NOT NULL, + internet_security_raw numeric NOT NULL, + stability_raw numeric NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT pk_index_raw_snapshot PRIMARY KEY (snapshot_date, member_id) +); + +CREATE INDEX IF NOT EXISTS idx_index_raw_snapshot_member + ON index_raw_snapshot (member_id); + +CREATE TABLE IF NOT EXISTS index_tscore_snapshot ( + snapshot_date date NOT NULL, + member_id bigint NOT NULL, + explore_tscore numeric NOT NULL, + benefit_trend_tscore numeric NOT NULL, + multi_device_tscore numeric NOT NULL, + family_home_tscore numeric NOT NULL, + internet_security_tscore numeric NOT NULL, + stability_tscore numeric NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT pk_index_tscore_snapshot PRIMARY KEY (snapshot_date, member_id) +); + +CREATE INDEX IF NOT EXISTS idx_index_tscore_snapshot_member + ON index_tscore_snapshot (member_id); + +CREATE TABLE IF NOT EXISTS persona_type ( + persona_type_id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, + character_name VARCHAR(100) NOT NULL, + short_desc TEXT, + character_description TEXT, + version INTEGER NOT NULL DEFAULT 1, + is_active BOOLEAN NOT NULL DEFAULT TRUE, + tags TEXT[] NOT NULL DEFAULT '{}', + created_at TIMESTAMP NOT NULL DEFAULT NOW(), + updated_at TIMESTAMP NOT NULL DEFAULT NOW(), + CONSTRAINT uk_persona_type_name_version UNIQUE (character_name, version) +); + +CREATE INDEX IF NOT EXISTS idx_persona_type_active_version + ON persona_type (is_active, version DESC); + +CREATE INDEX IF NOT EXISTS idx_persona_type_name + ON persona_type (character_name); + +CREATE TABLE IF NOT EXISTS index_persona_snapshot ( + snapshot_date date NOT NULL, + member_id bigint NOT NULL, + persona_type_id bigint NOT NULL, + -- 하위 호환 및 운영 조회 편의를 위해 코드도 함께 저장한다. + persona_code varchar(50) NOT NULL, + source_index_code varchar(50) NOT NULL, + source_tscore numeric NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT pk_index_persona_snapshot PRIMARY KEY (snapshot_date, member_id), + CONSTRAINT fk_index_persona_snapshot_member + FOREIGN KEY (member_id) REFERENCES member (member_id) + ON UPDATE CASCADE + ON DELETE CASCADE, + CONSTRAINT fk_index_persona_snapshot_persona_type + FOREIGN KEY (persona_type_id) REFERENCES persona_type (persona_type_id) +); + +CREATE INDEX IF NOT EXISTS idx_index_persona_snapshot_member + ON index_persona_snapshot (member_id); + +CREATE INDEX IF NOT EXISTS idx_index_persona_snapshot_persona + ON index_persona_snapshot (persona_code); + +CREATE INDEX IF NOT EXISTS idx_index_persona_snapshot_persona_type_id + ON index_persona_snapshot (persona_type_id); \ No newline at end of file diff --git a/src/main/resources/sql/index/delete_persona_snapshot.sql b/src/main/resources/sql/index/delete_persona_snapshot.sql new file mode 100644 index 0000000..5129b9d --- /dev/null +++ b/src/main/resources/sql/index/delete_persona_snapshot.sql @@ -0,0 +1,3 @@ +-- 목적: 동일 snapshot_date 재실행 시 persona 결과를 덮어쓰기 위해 기존 데이터를 삭제한다. +DELETE FROM index_persona_snapshot +WHERE snapshot_date = CAST(:snapshotDate AS date); diff --git a/src/main/resources/sql/index/delete_raw_snapshot.sql b/src/main/resources/sql/index/delete_raw_snapshot.sql new file mode 100644 index 0000000..e21c200 --- /dev/null +++ b/src/main/resources/sql/index/delete_raw_snapshot.sql @@ -0,0 +1,3 @@ +-- 목적: 동일 snapshot_date 재실행 시 raw 결과를 덮어쓰기 위해 기존 데이터를 삭제한다. +DELETE FROM index_raw_snapshot +WHERE snapshot_date = CAST(:snapshotDate AS date); diff --git a/src/main/resources/sql/index/delete_tscore_snapshot.sql b/src/main/resources/sql/index/delete_tscore_snapshot.sql new file mode 100644 index 0000000..ea9cbf7 --- /dev/null +++ b/src/main/resources/sql/index/delete_tscore_snapshot.sql @@ -0,0 +1,3 @@ +-- 목적: 동일 snapshot_date 재실행 시 tscore 결과를 덮어쓰기 위해 기존 데이터를 삭제한다. +DELETE FROM index_tscore_snapshot +WHERE snapshot_date = CAST(:snapshotDate AS date); diff --git a/src/main/resources/sql/index/insert_persona_snapshot.sql b/src/main/resources/sql/index/insert_persona_snapshot.sql new file mode 100644 index 0000000..916133c --- /dev/null +++ b/src/main/resources/sql/index/insert_persona_snapshot.sql @@ -0,0 +1,79 @@ +-- 목적: +-- 1) 회원별 6개 T-score를 세로로 펼친 뒤 최고 점수 지수 1개를 선택한다. +-- 2) 선택된 지수 코드를 persona character_name으로 변환한다. +-- 3) persona_type(활성 버전 우선)에서 persona_type_id를 찾아 FK로 저장한다. +-- +-- 동점 처리: +-- - tscore DESC, index_code ASC +-- - 즉, 점수가 같으면 index_code 사전순이 우선된다. +WITH expanded AS ( + SELECT + t.member_id, + v.index_code, + v.tscore + FROM index_tscore_snapshot t + CROSS JOIN LATERAL ( + VALUES + ('explore', t.explore_tscore), + ('benefit_trend', t.benefit_trend_tscore), + ('multi_device', t.multi_device_tscore), + ('family_home', t.family_home_tscore), + ('internet_security', t.internet_security_tscore), + ('stability', t.stability_tscore) + ) v(index_code, tscore) + WHERE t.snapshot_date = CAST(:snapshotDate AS date) +), +ranked AS ( + SELECT + e.member_id, + e.index_code, + e.tscore, + ROW_NUMBER() OVER ( + PARTITION BY e.member_id + ORDER BY e.tscore DESC, e.index_code ASC + ) AS rn + FROM expanded e +), +winners AS ( + SELECT + r.member_id, + r.index_code AS source_index_code, + r.tscore AS source_tscore, + CASE r.index_code + WHEN 'explore' THEN 'SPACE_SHERLOCK' + WHEN 'benefit_trend' THEN 'SPACE_SURFER' + WHEN 'multi_device' THEN 'SPACE_OCTOPUS' + WHEN 'family_home' THEN 'SPACE_GRAVITY' + WHEN 'internet_security' THEN 'SPACE_GUARDIAN' + WHEN 'stability' THEN 'SPACE_EXPLORER' + ELSE 'SPACE_SHERLOCK' + END AS persona_code + FROM ranked r + WHERE r.rn = 1 +), +latest_persona AS ( + SELECT DISTINCT ON (p.character_name) + p.character_name, + p.persona_type_id + FROM persona_type p + WHERE p.is_active = TRUE + ORDER BY p.character_name, p.version DESC, p.persona_type_id DESC +) +INSERT INTO index_persona_snapshot ( + snapshot_date, + member_id, + persona_type_id, + persona_code, + source_index_code, + source_tscore +) +SELECT + CAST(:snapshotDate AS date) AS snapshot_date, + w.member_id, + lp.persona_type_id, + w.persona_code, + w.source_index_code, + w.source_tscore +FROM winners w +JOIN latest_persona lp + ON lp.character_name = w.persona_code; \ No newline at end of file diff --git a/src/main/resources/sql/index/insert_raw_snapshot.sql b/src/main/resources/sql/index/insert_raw_snapshot.sql new file mode 100644 index 0000000..f361a31 --- /dev/null +++ b/src/main/resources/sql/index/insert_raw_snapshot.sql @@ -0,0 +1,248 @@ +-- 목적: +-- 1) 회원 전체를 모집단으로 6개 raw 지수를 계산한다. +-- 2) 로그/월 데이터 누락 회원도 포함하기 위해 LEFT JOIN + COALESCE 전략을 쓴다. +-- 3) 같은 snapshot_date 재실행 시 결과를 재현할 수 있도록 set-based 일괄 적재를 사용한다. +WITH params AS ( + -- 배치 기준 날짜와 월 키를 한 곳에서 고정한다. + SELECT + CAST(:snapshotDate AS date) AS snapshot_date, + CAST(:yyyymm AS varchar(6)) AS yyyymm +), +base_members AS ( + -- 모집단은 member 전체다. + SELECT + m.member_id, + m.membership, + COALESCE(m.children_count, 0) AS children_count, + m.family_group_id + FROM member m +), +family_counts AS ( + -- 가족 결합도 계산용: family_group_id별 인원 수 + SELECT + m.family_group_id, + COUNT(*)::bigint AS family_member_cnt + FROM member m + WHERE m.family_group_id IS NOT NULL + GROUP BY m.family_group_id +), +active_subscriptions AS ( + -- 스냅샷 시점에 유효한 활성 구독만 선택한다. + SELECT + s.member_id, + s.subscription_id, + s.product_id, + p.product_type, + p.name + FROM subscription s + JOIN product p ON p.product_id = s.product_id + JOIN params pr ON TRUE + WHERE s.status = TRUE + AND s.start_date::date <= pr.snapshot_date + AND (s.end_date IS NULL OR s.end_date::date >= pr.snapshot_date) +), +addon_features AS ( + -- 혜택/보안 지수 계산용 부가서비스 피처 + SELECT + a.member_id, + COUNT(*) FILTER (WHERE ad.addon_type::text = 'SECURITY')::bigint AS security_addon_cnt, + COUNT(*)::bigint AS addon_subscribe_cnt + FROM active_subscriptions a + JOIN addon_service ad ON ad.product_id = a.product_id + GROUP BY a.member_id +), +internet_iptv_features AS ( + -- 인터넷 가입 여부 + 인터넷/IPTV 동시 가입 여부 + SELECT + a.member_id, + CASE WHEN COUNT(*) FILTER (WHERE a.product_type::text = 'INTERNET') > 0 THEN 1 ELSE 0 END AS internet_flag, + CASE + WHEN COUNT(*) FILTER (WHERE a.product_type::text = 'INTERNET') > 0 + AND COUNT(*) FILTER (WHERE a.product_type::text = 'IPTV') > 0 + THEN 1 ELSE 0 + END AS has_internet_iptv_bundle + FROM active_subscriptions a + GROUP BY a.member_id +), +watch_tablet_features AS ( + -- TAB_WATCH_PLAN 중 이름 패턴으로 워치/태블릿을 구분한다. + SELECT + a.member_id, + CASE + WHEN COUNT(*) FILTER ( + WHERE a.product_type::text = 'TAB_WATCH_PLAN' + AND ( + a.name ILIKE '%watch%' + OR a.name ILIKE '%wearable%' + ) + ) > 0 THEN 1 ELSE 0 + END AS watch_flag, + CASE + WHEN COUNT(*) FILTER ( + WHERE a.product_type::text = 'TAB_WATCH_PLAN' + AND NOT ( + a.name ILIKE '%watch%' + OR a.name ILIKE '%wearable%' + ) + ) > 0 THEN 1 ELSE 0 + END AS tablet_flag + FROM active_subscriptions a + GROUP BY a.member_id +), +mobile_limit AS ( + -- 제공량 문자열에서 숫자만 추출해 GB 한도로 파싱한다. + -- 이유: 실제 데이터가 "기본제공량 내 ... 55GB" 같은 문자열 포맷을 포함함. + SELECT + a.member_id, + COALESCE( + MAX( + CASE + WHEN mp.tethering_sharing_data IS NULL THEN NULL + WHEN regexp_replace(mp.tethering_sharing_data::text, '[^0-9.]', '', 'g') + ~ '^[0-9]+(\\.[0-9]+)?$' + THEN regexp_replace(mp.tethering_sharing_data::text, '[^0-9.]', '', 'g')::numeric + ELSE NULL + END + ), + 0 + ) AS sharing_limit_gb + FROM active_subscriptions a + JOIN mobile_plan mp ON mp.product_id = a.product_id + WHERE a.product_type::text = 'MOBILE_PLAN' + GROUP BY a.member_id +), +sharing_used AS ( + -- usage JSON에서 tethering_sharing_data_gb를 숫자형으로 안전 추출한다. + SELECT + a.member_id, + SUM( + CASE + WHEN jsonb_exists(um.usage_details, 'tethering_sharing_data_gb') + AND (um.usage_details ->> 'tethering_sharing_data_gb') ~ '^[0-9]+(\\.[0-9]+)?$' + THEN (um.usage_details ->> 'tethering_sharing_data_gb')::numeric + ELSE 0 + END + ) AS sharing_used_gb + FROM active_subscriptions a + JOIN usage_monthly um ON um.subscription_id = a.subscription_id + JOIN params pr ON TRUE + WHERE a.product_type::text = 'MOBILE_PLAN' + AND um.yyyymm = pr.yyyymm + GROUP BY a.member_id +), +coupon_7d AS ( + -- 스냅샷 기준 최근 7일 쿠폰 사용 건수 + SELECT + mc.member_id, + COUNT(*)::bigint AS coupon_used_cnt + FROM member_coupon mc + JOIN params pr ON TRUE + WHERE mc.is_used = TRUE + AND mc.used_at IS NOT NULL + AND mc.used_at::date BETWEEN (pr.snapshot_date - INTERVAL '6 day')::date AND pr.snapshot_date + GROUP BY mc.member_id +), +support_monthly AS ( + -- 스냅샷 월(yyyymm) 상담 만족도 평균 + SELECT + sc.member_id, + AVG(sc.satisfaction_score::numeric) AS avg_monthly_satisfaction + FROM support_case sc + JOIN params pr ON TRUE + WHERE sc.satisfaction_score IS NOT NULL + AND to_char(COALESCE(sc.resolved_at, sc.updated_at), 'YYYYMM') = pr.yyyymm + GROUP BY sc.member_id +), +billing_monthly AS ( + -- 스냅샷 월 납부 상태 + SELECT + b.member_id, + b.is_paid + FROM billing b + JOIN params pr ON TRUE + WHERE b.yyyymm = pr.yyyymm +), +log_features AS ( + -- Athena 집계 적재본에서 탐색/액션 관련 로그 피처를 가져온다. + SELECT + l.member_id, + l.click_product_detail_cnt, + l.click_compare_cnt, + l.click_change_success_cnt + FROM user_event_features_7d l + JOIN params pr ON l.snapshot_date = pr.snapshot_date +) +INSERT INTO index_raw_snapshot ( + snapshot_date, + member_id, + explore_raw, + benefit_trend_raw, + multi_device_raw, + family_home_raw, + internet_security_raw, + stability_raw +) +SELECT + pr.snapshot_date AS snapshot_date, + bm.member_id, + -- 탐색 지수: 상세조회 + 비교(가중치 3) + ( + LN(1 + COALESCE(lf.click_product_detail_cnt, 0)::numeric) + + 3 * LN(1 + COALESCE(lf.click_compare_cnt, 0)::numeric) + ) AS explore_raw, + -- 혜택/트렌드: 부가가입 + 쿠폰사용 + 성공액션 + ( + LN(1 + COALESCE(af.addon_subscribe_cnt, 0)::numeric) + + LN(1 + COALESCE(c7.coupon_used_cnt, 0)::numeric) + + LN(1 + COALESCE(lf.click_change_success_cnt, 0)::numeric) + ) AS benefit_trend_raw, + -- 멀티 디바이스: 워치/태블릿 플래그 + sharing_rate(0~1 clamp) + ( + COALESCE(wtf.watch_flag, 0) + + COALESCE(wtf.tablet_flag, 0) + + CASE + WHEN COALESCE(ml.sharing_limit_gb, 0) > 0 + THEN LEAST( + 1.0::numeric, + GREATEST(0.0::numeric, COALESCE(su.sharing_used_gb, 0) / ml.sharing_limit_gb) + ) + ELSE 0.0::numeric + END + ) AS multi_device_raw, + -- 가족/홈: 가족인원 + 홈결합 + 자녀수 + ( + LN(1 + COALESCE(fc.family_member_cnt, 0)::numeric) + + COALESCE(iif.has_internet_iptv_bundle, 0) + + LN(1 + COALESCE(bm.children_count, 0)::numeric) + ) AS family_home_raw, + -- 인터넷/보안: 인터넷 가입 신호(0.7) + 보안 부가서비스 신호(0.3) + ( + 0.7 * COALESCE(iif.internet_flag, 0)::numeric + + 0.3 * LN(1 + COALESCE(af.security_addon_cnt, 0)::numeric) + ) AS internet_security_raw, + -- 안정성: 멤버십/상담평점/납부상태 평균 + ( + ( + CASE bm.membership + WHEN 'VVIP' THEN 100 + WHEN 'VIP' THEN 80 + WHEN 'GOLD' THEN 60 + WHEN 'BASIC' THEN 40 + ELSE 0 + END + + (COALESCE(sm.avg_monthly_satisfaction, 0) / 5.0) * 100 + + CASE WHEN COALESCE(bl.is_paid, FALSE) THEN 100 ELSE 0 END + ) / 3.0 + ) AS stability_raw +FROM base_members bm +JOIN params pr ON TRUE +LEFT JOIN family_counts fc ON fc.family_group_id = bm.family_group_id +LEFT JOIN addon_features af ON af.member_id = bm.member_id +LEFT JOIN internet_iptv_features iif ON iif.member_id = bm.member_id +LEFT JOIN watch_tablet_features wtf ON wtf.member_id = bm.member_id +LEFT JOIN mobile_limit ml ON ml.member_id = bm.member_id +LEFT JOIN sharing_used su ON su.member_id = bm.member_id +LEFT JOIN coupon_7d c7 ON c7.member_id = bm.member_id +LEFT JOIN support_monthly sm ON sm.member_id = bm.member_id +LEFT JOIN billing_monthly bl ON bl.member_id = bm.member_id +LEFT JOIN log_features lf ON lf.member_id = bm.member_id; diff --git a/src/main/resources/sql/index/insert_tscore_snapshot.sql b/src/main/resources/sql/index/insert_tscore_snapshot.sql new file mode 100644 index 0000000..7f131b8 --- /dev/null +++ b/src/main/resources/sql/index/insert_tscore_snapshot.sql @@ -0,0 +1,74 @@ +-- 목적: +-- 1) raw 스냅샷에서 지수별 평균/표준편차를 계산한다. +-- 2) T = 50 + 10 * Z 공식을 적용해 tscore를 적재한다. +-- 3) 표준편차가 0인 경우 분모 0 방지를 위해 50점으로 고정한다. +WITH raw_data AS ( + -- 특정 snapshot_date 데이터만 대상으로 한다. + SELECT + r.member_id, + r.explore_raw, + r.benefit_trend_raw, + r.multi_device_raw, + r.family_home_raw, + r.internet_security_raw, + r.stability_raw + FROM index_raw_snapshot r + WHERE r.snapshot_date = CAST(:snapshotDate AS date) +), +stats AS ( + -- 지수별 분포 통계량을 한 번에 계산한다. + SELECT + AVG(explore_raw) AS explore_mean, + STDDEV_POP(explore_raw) AS explore_stddev, + AVG(benefit_trend_raw) AS benefit_trend_mean, + STDDEV_POP(benefit_trend_raw) AS benefit_trend_stddev, + AVG(multi_device_raw) AS multi_device_mean, + STDDEV_POP(multi_device_raw) AS multi_device_stddev, + AVG(family_home_raw) AS family_home_mean, + STDDEV_POP(family_home_raw) AS family_home_stddev, + AVG(internet_security_raw) AS internet_security_mean, + STDDEV_POP(internet_security_raw) AS internet_security_stddev, + AVG(stability_raw) AS stability_mean, + STDDEV_POP(stability_raw) AS stability_stddev + FROM raw_data +) +INSERT INTO index_tscore_snapshot ( + snapshot_date, + member_id, + explore_tscore, + benefit_trend_tscore, + multi_device_tscore, + family_home_tscore, + internet_security_tscore, + stability_tscore +) +SELECT + CAST(:snapshotDate AS date) AS snapshot_date, + r.member_id, + -- stddev=0이면 모두 같은 값이므로 중립점수 50 부여 + CASE + WHEN COALESCE(s.explore_stddev, 0) = 0 THEN 50 + ELSE 50 + 10 * ((r.explore_raw - s.explore_mean) / s.explore_stddev) + END AS explore_tscore, + CASE + WHEN COALESCE(s.benefit_trend_stddev, 0) = 0 THEN 50 + ELSE 50 + 10 * ((r.benefit_trend_raw - s.benefit_trend_mean) / s.benefit_trend_stddev) + END AS benefit_trend_tscore, + CASE + WHEN COALESCE(s.multi_device_stddev, 0) = 0 THEN 50 + ELSE 50 + 10 * ((r.multi_device_raw - s.multi_device_mean) / s.multi_device_stddev) + END AS multi_device_tscore, + CASE + WHEN COALESCE(s.family_home_stddev, 0) = 0 THEN 50 + ELSE 50 + 10 * ((r.family_home_raw - s.family_home_mean) / s.family_home_stddev) + END AS family_home_tscore, + CASE + WHEN COALESCE(s.internet_security_stddev, 0) = 0 THEN 50 + ELSE 50 + 10 * ((r.internet_security_raw - s.internet_security_mean) / s.internet_security_stddev) + END AS internet_security_tscore, + CASE + WHEN COALESCE(s.stability_stddev, 0) = 0 THEN 50 + ELSE 50 + 10 * ((r.stability_raw - s.stability_mean) / s.stability_stddev) + END AS stability_tscore +FROM raw_data r +CROSS JOIN stats s; diff --git a/src/main/resources/sql/index/verify_counts.sql b/src/main/resources/sql/index/verify_counts.sql new file mode 100644 index 0000000..b764593 --- /dev/null +++ b/src/main/resources/sql/index/verify_counts.sql @@ -0,0 +1,6 @@ +-- 목적: 모집단(member) 대비 raw/tscore/persona 누락 여부를 건수로 검증한다. +SELECT + (SELECT COUNT(*) FROM member) AS member_count, + (SELECT COUNT(*) FROM index_raw_snapshot WHERE snapshot_date = CAST(:snapshotDate AS date)) AS raw_count, + (SELECT COUNT(*) FROM index_tscore_snapshot WHERE snapshot_date = CAST(:snapshotDate AS date)) AS tscore_count, + (SELECT COUNT(*) FROM index_persona_snapshot WHERE snapshot_date = CAST(:snapshotDate AS date)) AS persona_count; diff --git a/src/main/resources/sql/index/verify_persona_nulls.sql b/src/main/resources/sql/index/verify_persona_nulls.sql new file mode 100644 index 0000000..0e1eb31 --- /dev/null +++ b/src/main/resources/sql/index/verify_persona_nulls.sql @@ -0,0 +1,9 @@ +-- 목적: persona 스냅샷 필수 컬럼 null 여부 검증 +SELECT COUNT(*) AS null_count +FROM index_persona_snapshot +WHERE snapshot_date = CAST(:snapshotDate AS date) + AND ( + persona_type_id IS NULL + OR source_index_code IS NULL + OR source_tscore IS NULL + ); \ No newline at end of file diff --git a/src/main/resources/sql/index/verify_tscore_nulls.sql b/src/main/resources/sql/index/verify_tscore_nulls.sql new file mode 100644 index 0000000..a3307ef --- /dev/null +++ b/src/main/resources/sql/index/verify_tscore_nulls.sql @@ -0,0 +1,12 @@ +-- 목적: tscore 필수 컬럼에 null 값이 있는지 확인한다. +SELECT COUNT(*) AS null_count +FROM index_tscore_snapshot +WHERE snapshot_date = CAST(:snapshotDate AS date) + AND ( + explore_tscore IS NULL + OR benefit_trend_tscore IS NULL + OR multi_device_tscore IS NULL + OR family_home_tscore IS NULL + OR internet_security_tscore IS NULL + OR stability_tscore IS NULL + ); From 48ef86ca59909d3cbcc8cd3a124bff63d44a1f72 Mon Sep 17 00:00:00 2001 From: bon0512 Date: Thu, 12 Mar 2026 13:02:35 +0900 Subject: [PATCH 5/5] =?UTF-8?q?[HSC-200]=20feat:=20=EC=A3=BC=EC=84=9D,=20?= =?UTF-8?q?=ED=85=8C=EC=9D=B4=EB=B8=94=20=EC=83=9D=EC=84=B1sql=20=EB=AC=B8?= =?UTF-8?q?=20=EC=82=AD=EC=A0=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../BuildPersonaWeeklyIndexTasklet.java | 4 +- .../tasklet/BuildRawWeeklyIndexTasklet.java | 4 +- .../BuildTScoreWeeklyIndexTasklet.java | 4 +- .../index/tasklet/GateWeeklyIndexTasklet.java | 2 +- .../tasklet/VerifyWeeklyIndexTasklet.java | 4 +- .../index/create_index_snapshot_tables.sql | 82 ------------------- 6 files changed, 9 insertions(+), 91 deletions(-) delete mode 100644 src/main/resources/sql/index/create_index_snapshot_tables.sql diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java index edee6da..189777d 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildPersonaWeeklyIndexTasklet.java @@ -55,8 +55,8 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon int deleted = namedParameterJdbcTemplate.update(deleteSql, params); int inserted = namedParameterJdbcTemplate.update(insertSql, params); - log.info("BuildPersona 완료. snapshotDate={}, deleted={}, inserted={}", - snapshotDate, deleted, inserted); + /*log.info("BuildPersona 완료. snapshotDate={}, deleted={}, inserted={}", + snapshotDate, deleted, inserted);*/ return RepeatStatus.FINISHED; } diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java index 35c040e..bbef84f 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildRawWeeklyIndexTasklet.java @@ -60,8 +60,8 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon int deleted = namedParameterJdbcTemplate.update(deleteSql, params); int inserted = namedParameterJdbcTemplate.update(insertSql, params); - log.info("BuildRaw 완료. snapshotDate={}, yyyymm={}, deleted={}, inserted={}", - snapshotDate, yyyymm, deleted, inserted); + /*log.info("BuildRaw 완료. snapshotDate={}, yyyymm={}, deleted={}, inserted={}", + snapshotDate, yyyymm, deleted, inserted);*/ return RepeatStatus.FINISHED; } diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java index fc5e87d..35cce64 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/BuildTScoreWeeklyIndexTasklet.java @@ -50,8 +50,8 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon int deleted = namedParameterJdbcTemplate.update(deleteSql, params); int inserted = namedParameterJdbcTemplate.update(insertSql, params); - log.info("BuildTScore 완료. snapshotDate={}, deleted={}, inserted={}", - snapshotDate, deleted, inserted); + /*log.info("BuildTScore 완료. snapshotDate={}, deleted={}, inserted={}", + snapshotDate, deleted, inserted);*/ return RepeatStatus.FINISHED; } diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java index c6c256f..f69f905 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/GateWeeklyIndexTasklet.java @@ -82,7 +82,7 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon // 4) 필수 테이블 존재 여부를 확인한다. verifyRequiredTables(); - log.info("Gate 통과. snapshotDate={}, yyyymm={}", snapshotDate, yyyymm); + //log.info("Gate 통과. snapshotDate={}, yyyymm={}", snapshotDate, yyyymm); return RepeatStatus.FINISHED; } diff --git a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java index dcb7c35..0608fd1 100644 --- a/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java +++ b/src/main/java/site/holliverse/worker/batch/jobs/index/tasklet/VerifyWeeklyIndexTasklet.java @@ -77,7 +77,7 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon throw new IllegalStateException("Persona null 검증에 실패했습니다. nullCount=" + personaNullCount); } - log.info( + /*log.info( "Verify 완료. snapshotDate={}, memberCount={}, rawCount={}, tscoreCount={}, personaCount={}, tscoreNullCount={}, personaNullCount={}", snapshotDate, memberCount, @@ -86,7 +86,7 @@ public RepeatStatus execute(StepContribution contribution, ChunkContext chunkCon personaCount, tscoreNullCount, personaNullCount - ); + );*/ return RepeatStatus.FINISHED; } diff --git a/src/main/resources/sql/index/create_index_snapshot_tables.sql b/src/main/resources/sql/index/create_index_snapshot_tables.sql deleted file mode 100644 index 407de6a..0000000 --- a/src/main/resources/sql/index/create_index_snapshot_tables.sql +++ /dev/null @@ -1,82 +0,0 @@ --- 목적: 주간 지수 배치 결과(raw/tscore/persona) 저장 테이블 생성 --- 특이사항: --- - snapshot_date + member_id 를 기본 키로 사용해 동일 날짜 중복 적재를 방지한다. --- - persona 결과는 persona_type_id(FK) 기준으로 저장한다. - -CREATE TABLE IF NOT EXISTS index_raw_snapshot ( - snapshot_date date NOT NULL, - member_id bigint NOT NULL, - explore_raw numeric NOT NULL, - benefit_trend_raw numeric NOT NULL, - multi_device_raw numeric NOT NULL, - family_home_raw numeric NOT NULL, - internet_security_raw numeric NOT NULL, - stability_raw numeric NOT NULL, - created_at timestamptz NOT NULL DEFAULT now(), - CONSTRAINT pk_index_raw_snapshot PRIMARY KEY (snapshot_date, member_id) -); - -CREATE INDEX IF NOT EXISTS idx_index_raw_snapshot_member - ON index_raw_snapshot (member_id); - -CREATE TABLE IF NOT EXISTS index_tscore_snapshot ( - snapshot_date date NOT NULL, - member_id bigint NOT NULL, - explore_tscore numeric NOT NULL, - benefit_trend_tscore numeric NOT NULL, - multi_device_tscore numeric NOT NULL, - family_home_tscore numeric NOT NULL, - internet_security_tscore numeric NOT NULL, - stability_tscore numeric NOT NULL, - created_at timestamptz NOT NULL DEFAULT now(), - CONSTRAINT pk_index_tscore_snapshot PRIMARY KEY (snapshot_date, member_id) -); - -CREATE INDEX IF NOT EXISTS idx_index_tscore_snapshot_member - ON index_tscore_snapshot (member_id); - -CREATE TABLE IF NOT EXISTS persona_type ( - persona_type_id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, - character_name VARCHAR(100) NOT NULL, - short_desc TEXT, - character_description TEXT, - version INTEGER NOT NULL DEFAULT 1, - is_active BOOLEAN NOT NULL DEFAULT TRUE, - tags TEXT[] NOT NULL DEFAULT '{}', - created_at TIMESTAMP NOT NULL DEFAULT NOW(), - updated_at TIMESTAMP NOT NULL DEFAULT NOW(), - CONSTRAINT uk_persona_type_name_version UNIQUE (character_name, version) -); - -CREATE INDEX IF NOT EXISTS idx_persona_type_active_version - ON persona_type (is_active, version DESC); - -CREATE INDEX IF NOT EXISTS idx_persona_type_name - ON persona_type (character_name); - -CREATE TABLE IF NOT EXISTS index_persona_snapshot ( - snapshot_date date NOT NULL, - member_id bigint NOT NULL, - persona_type_id bigint NOT NULL, - -- 하위 호환 및 운영 조회 편의를 위해 코드도 함께 저장한다. - persona_code varchar(50) NOT NULL, - source_index_code varchar(50) NOT NULL, - source_tscore numeric NOT NULL, - created_at timestamptz NOT NULL DEFAULT now(), - CONSTRAINT pk_index_persona_snapshot PRIMARY KEY (snapshot_date, member_id), - CONSTRAINT fk_index_persona_snapshot_member - FOREIGN KEY (member_id) REFERENCES member (member_id) - ON UPDATE CASCADE - ON DELETE CASCADE, - CONSTRAINT fk_index_persona_snapshot_persona_type - FOREIGN KEY (persona_type_id) REFERENCES persona_type (persona_type_id) -); - -CREATE INDEX IF NOT EXISTS idx_index_persona_snapshot_member - ON index_persona_snapshot (member_id); - -CREATE INDEX IF NOT EXISTS idx_index_persona_snapshot_persona - ON index_persona_snapshot (persona_code); - -CREATE INDEX IF NOT EXISTS idx_index_persona_snapshot_persona_type_id - ON index_persona_snapshot (persona_type_id); \ No newline at end of file