Commit e4f6b39d authored by sheteng's avatar sheteng

订单表写入实时汇率

parent 98f00dec
......@@ -21,7 +21,7 @@ version = '0.3'
mainClassName = 'com.fshark.overseas.advert.AppsflyerAnalyze'
ext {
javaVersion = '1.8'
flinkVersion = '1.12.0'
flinkVersion = '1.11.1'
scalaBinaryVersion = '2.12'
slf4jVersion = '1.7.7'
log4jVersion = '1.2.17'
......@@ -55,10 +55,10 @@ dependencies {
// https://mvnrepository.com/artifact/mysql/mysql-connector-java
implementation group: 'mysql', name: 'mysql-connector-java', version: '8.0.25'
compileOnly "org.apache.flink:flink-statebackend-rocksdb_${scalaBinaryVersion}:${flinkVersion}"
compileOnly "org.apache.flink:flink-java:${flinkVersion}"
compileOnly "org.apache.flink:flink-streaming-java_${scalaBinaryVersion}:${flinkVersion}"
compileOnly "org.apache.flink:flink-clients_${scalaBinaryVersion}:${flinkVersion}"
compile "org.apache.flink:flink-statebackend-rocksdb_${scalaBinaryVersion}:${flinkVersion}"
compile "org.apache.flink:flink-java:${flinkVersion}"
compile "org.apache.flink:flink-streaming-java_${scalaBinaryVersion}:${flinkVersion}"
compile "org.apache.flink:flink-clients_${scalaBinaryVersion}:${flinkVersion}"
implementation 'ru.yandex.clickhouse:clickhouse-jdbc:0.1.52'
// --------------------------------------------------------------
// Dependencies that should be part of the shadow jar, e.g.
......
......@@ -136,7 +136,7 @@ public class AppsflyerAnalyze {
DataStream<AppsFlyerEvent> sdkRegisterStream = mainStream.getSideOutput(sdkRegisterTag);
// loginStream.print();
orderStream.print();
//广告指标(激活)
DataStream<ActivationMetrics> activitionMetricStream = installStream
.flatMap(new ActivitionProcess(appConfig))
......
......@@ -174,6 +174,6 @@ public class AppsflyerMetricStore {
.uid("event-sink-event")
.name("event-sink-event");
sse.execute("AppsflyerMetricStore-analyze");
sse.execute("Appsflyer-Store-analyze");
}
}
......@@ -2,10 +2,12 @@ package com.fshark.overseas.advert.entity.dao;
import com.fshark.overseas.advert.entity.mapper.AdMappingMapper;
import com.fshark.overseas.advert.entity.mapper.GameMapper;
import com.fshark.overseas.advert.entity.mapper.MediaCampaignMapper;
import com.fshark.overseas.advert.entity.mapper.MediaMapper;
import com.fshark.overseas.advert.modle.AppsFlyerEvent;
import com.fshark.overseas.advert.modle.mapping.AdMapping;
import com.fshark.overseas.advert.modle.mapping.Game;
import com.fshark.overseas.advert.modle.mapping.Media;
import com.fshark.overseas.advert.modle.mapping.MediaCampaign;
import org.jdbi.v3.sqlobject.config.RegisterRowMapper;
......@@ -23,8 +25,9 @@ public interface GameMappingDao {
/**
* 拉取游戏
*/
@SqlQuery("select count(1) from game where platform_id = :platform_id and af_app_id = :game_id")
boolean findGameByPrimary(@BindBean AppsFlyerEvent event, @Bind("platform_id") int platform_id);
@SqlQuery("select * from game where af_app_id = :game_id and platform_id = :platform_id ")
@RegisterRowMapper(GameMapper.class)
Optional<Game> findGameByPrimary(@Bind("game_id") String game_id,@Bind("platform_id")int platform_id);
@SqlUpdate("insert into game (`platform_id`, `af_app_id`, `af_app_name`, `app_id`,`select_time_zone`) values (:platform_id,:game_id,:app_name,:app_id,:select_time_zone)")
......
......@@ -8,9 +8,9 @@ public interface OrderDao {
@SqlUpdate("insert into appsflyer_pay_record (region, did, uid,platform, game_id, media_source, campaign_id, ad_set_id,ad_id, " +
" pay_time, money, currency, order_no, reg_time, activation_time, time_zone ,channel,data_type,retargeting_time)" +
" pay_time, money, currency, order_no, reg_time, activation_time, time_zone ,channel,data_type,retargeting_time,rate)" +
" values(:region, :did, :uid, :platform, :game_id, :media_source, :campaign_id, :ad_set_id , :ad_id," +
" :pay_time, :revenue,:currency,:orderno,:reg_time,:activation_time,:time_zone,:channel,:data_type,:retargeting_time)")
" :pay_time, :revenue,:currency,:orderno,:reg_time,:activation_time,:time_zone,:channel,:data_type,:retargeting_time,:rate)")
void insertOrder(@BindBean OrderData orderData);
// @SqlUpdate("insert into appsflyer_pay_record_server (region, did, uid,platform, game_id, media_source, campaign_id, ad_set_id,ad_id, " +
......@@ -20,9 +20,9 @@ public interface OrderDao {
// void insertServerOrder(@BindBean OrderData eventData);
@SqlUpdate("insert into appsflyer_pay_record_ouid (region, did, uid,ouid,platform, game_id, media_source, campaign_id, ad_set_id,ad_id, " +
" pay_time, money, currency, order_no, ouid_reg_time, time_zone ,channel,data_type)" +
" pay_time, money, currency, order_no, ouid_reg_time, time_zone ,channel,data_type,rate)" +
" values(:region, :did, :uid, :ouid,:platform, :game_id, :media_source, :campaign_id, :ad_set_id , :ad_id," +
" :pay_time, :revenue,:currency,:orderno,:ouid_reg_time,:time_zone,:channel,:data_type)")
" :pay_time, :revenue,:currency,:orderno,:ouid_reg_time,:time_zone,:channel,:data_type,:rate)")
void insertOuidOrder(@BindBean OrderData orderData);
// @SqlUpdate("insert into appsflyer_order_accumulative" +
......
package com.fshark.overseas.advert.entity.mapper;
import com.fshark.overseas.advert.modle.mapping.Game;
import com.fshark.overseas.advert.modle.mapping.Media;
import org.jdbi.v3.core.mapper.RowMapper;
import org.jdbi.v3.core.statement.StatementContext;
import java.sql.ResultSet;
import java.sql.SQLException;
public class GameMapper implements RowMapper<Game> {
@Override
public Game map(ResultSet rs, StatementContext ctx) throws SQLException {
return Game.builder()
.afAppId(rs.getString("af_app_id"))
.rate(rs.getDouble("rate"))
.build();
}
}
......@@ -154,16 +154,10 @@ public class OrderData implements Serializable {
// private String server_reg_hour;
private String ouid_reg_time;
private String ouid_platform;
private String ouid_media_source;
private String ouid_campaign_id;
private String ouid_ad_set_id;
private String ouid_ad_id;
/**
* 汇率
*/
private Double rate;
private int data_type;
......
......@@ -7,7 +7,8 @@ package com.fshark.overseas.advert.modle;
public enum PlatformEnum {
IOS("ios"),
Android("android");
Android("android"),
PC("pc");
PlatformEnum(String name) {
this.name = name;
......
package com.fshark.overseas.advert.modle.mapping;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
@Builder
@Data
@AllArgsConstructor
public class Game {
private Integer platformId;
private String afAppId;
private double rate;
}
\ No newline at end of file
......@@ -51,10 +51,10 @@ public class GameAdMappingProcessor extends RichFlatMapFunction<AppsFlyerEvent,
}
//更新game
boolean gameByPrimary = gameMappingDao.findGameByPrimary(value, platform_id);
if (!gameByPrimary) {
gameMappingDao.saveGame(value, platform_id);
}
// boolean gameByPrimary = gameMappingDao.findGameByPrimary(value, platform_id);
// if (!gameByPrimary) {
// gameMappingDao.saveGame(value, platform_id);
// }
//mediaSource更新
Optional<Media> mediaByPrimary = gameMappingDao.findMediaByPrimary(value);
......
......@@ -4,12 +4,14 @@ import com.alibaba.fastjson.JSONObject;
import com.fshark.overseas.advert.entity.GameAccount;
import com.fshark.overseas.advert.entity.ServerRoleAccount;
import com.fshark.overseas.advert.entity.dao.GameAccountDao;
import com.fshark.overseas.advert.entity.dao.GameMappingDao;
import com.fshark.overseas.advert.entity.dao.OrderDao;
import com.fshark.overseas.advert.entity.dao.ServerRoleDao;
import com.fshark.overseas.advert.modle.AppsFlyerEvent;
import com.fshark.overseas.advert.modle.AppsflyerPurchase;
import com.fshark.overseas.advert.modle.OrderData;
import com.fshark.overseas.advert.modle.SinkMetricEvent;
import com.fshark.overseas.advert.modle.mapping.Game;
import com.fshark.overseas.advert.modle.metrics.ActivationMetrics;
import com.fshark.overseas.advert.util.*;
import com.google.gson.Gson;
......@@ -29,7 +31,6 @@ public class OuidOrderDataProcessor extends RichFlatMapFunction<AppsFlyerEvent,
private JdbiContext jdbiContext;
private GameAccountDao gameAccountDao;
private OrderDao orderDao;
public OuidOrderDataProcessor(AppConfig appConfig) {
this.appConfig = appConfig;
......@@ -40,7 +41,6 @@ public class OuidOrderDataProcessor extends RichFlatMapFunction<AppsFlyerEvent,
super.open(parameters);
jdbiContext = DatabaseUtils.getGenJdbiClient(appConfig.getDbHost(), appConfig.getDbUser(), appConfig.getDbPassword());
gameAccountDao = jdbiContext.getJdbi().onDemand(GameAccountDao.class);
orderDao = jdbiContext.getJdbi().onDemand(OrderDao.class);
}
......@@ -78,6 +78,7 @@ public class OuidOrderDataProcessor extends RichFlatMapFunction<AppsFlyerEvent,
orderData.setRegion(gameAccount.getRegion());
orderData.setData_type(CommonUtils.setType(Long.parseLong(DateUtils.formatAsYMD(gameAccount.getEvent_time())), orderData.getMedia_source(), orderData.getData_type()));
}
SinkMetricEvent sinkMetricEvent = new SinkMetricEvent();
sinkMetricEvent.setEventType(Constants.TYPE_SINK_OUID_ORDER);
sinkMetricEvent.setOrderData(orderData);
......
......@@ -3,11 +3,13 @@ package com.fshark.overseas.advert.processor;
import com.fshark.overseas.advert.entity.GameAccount;
import com.fshark.overseas.advert.entity.ServerRoleAccount;
import com.fshark.overseas.advert.entity.dao.GameAccountDao;
import com.fshark.overseas.advert.entity.dao.GameMappingDao;
import com.fshark.overseas.advert.entity.dao.ServerRoleDao;
import com.fshark.overseas.advert.modle.AppsFlyerEvent;
import com.fshark.overseas.advert.modle.AppsflyerPurchase;
import com.fshark.overseas.advert.modle.OrderData;
import com.fshark.overseas.advert.modle.SinkMetricEvent;
import com.fshark.overseas.advert.modle.mapping.Game;
import com.fshark.overseas.advert.util.*;
import com.google.gson.Gson;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
......@@ -22,7 +24,7 @@ public class PurchaseToOrderDataProcessor extends RichFlatMapFunction<AppsFlyerE
private JdbiContext jdbiContext;
private GameAccountDao gameAccountDao;
private ServerRoleDao serverRoleDao;
private GameMappingDao gameMappingDao;
public PurchaseToOrderDataProcessor(AppConfig appConfig) {
this.appConfig = appConfig;
......@@ -33,7 +35,7 @@ public class PurchaseToOrderDataProcessor extends RichFlatMapFunction<AppsFlyerE
super.open(parameters);
jdbiContext = DatabaseUtils.getGenJdbiClient(appConfig.getDbHost(), appConfig.getDbUser(), appConfig.getDbPassword());
gameAccountDao = jdbiContext.getJdbi().onDemand(GameAccountDao.class);
// serverRoleDao = jdbiContext.getJdbi().onDemand(ServerRoleDao.class);
}
@Override
......
package com.fshark.overseas.advert.sink;
import com.fshark.overseas.advert.entity.dao.GameMappingDao;
import com.fshark.overseas.advert.entity.dao.OrderDao;
import com.fshark.overseas.advert.modle.OrderData;
import com.fshark.overseas.advert.modle.PlatformEnum;
import com.fshark.overseas.advert.modle.SinkMetricEvent;
import com.fshark.overseas.advert.modle.mapping.Game;
import com.fshark.overseas.advert.util.*;
import org.apache.commons.lang3.StringUtils;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import java.util.Optional;
public class PurchaseSink extends RichSinkFunction<SinkMetricEvent> {
private AppConfig appConfig;
private JdbiContext jdbiContext;
private OrderDao orderDao;
private GameMappingDao gameMappingDao;
public PurchaseSink(AppConfig appConfig) {
this.appConfig = appConfig;
......@@ -28,16 +34,37 @@ public class PurchaseSink extends RichSinkFunction<SinkMetricEvent> {
jdbiContext = DatabaseUtils.getCHGenJdbiClient(appConfig.getDbHost(), appConfig.getDbUser(), appConfig.getDbPassword());
}
orderDao = jdbiContext.getJdbi().onDemand(OrderDao.class);
gameMappingDao = jdbiContext.getJdbi().onDemand(GameMappingDao.class);
}
@Override
public void invoke(SinkMetricEvent value, Context context) throws Exception {
super.invoke(value, context);
if (Constants.TYPE_SINK_ORDER.equals(value.getEventType())) {
OrderData eventData = value.getOrderData();
//获取汇率
int platform_id = 1;
if (PlatformEnum.IOS.getName().equals(eventData.getPlatform())) {
platform_id = 2;
} else if (PlatformEnum.PC.getName().equals(eventData.getPlatform())){
platform_id = 3;
}
Optional<Game> gameByPrimary = gameMappingDao.findGameByPrimary(eventData.getGame_id(),platform_id);
gameByPrimary.ifPresent(game -> eventData.setRate(game.getRate()));
orderDao.insertOrder(eventData);
} else if (Constants.TYPE_SINK_OUID_ORDER.equals(value.getEventType())){
OrderData eventData = value.getOrderData();
//获取汇率
//获取汇率
int platform_id = 1;
if (PlatformEnum.IOS.getName().equals(eventData.getPlatform())) {
platform_id = 2;
} else if (PlatformEnum.PC.getName().equals(eventData.getPlatform())){
platform_id = 3;
}
Optional<Game> gameByPrimary = gameMappingDao.findGameByPrimary(eventData.getGame_id(),platform_id);
gameByPrimary.ifPresent(game -> eventData.setRate(game.getRate()));
orderDao.insertOuidOrder(eventData);
}
}
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment