diff --git a/config_templates/sql/initDB.sql b/config_templates/sql/initDB.sql index 09b463d..0bed689 100644 --- a/config_templates/sql/initDB.sql +++ b/config_templates/sql/initDB.sql @@ -1,16 +1,26 @@ +DROP TABLE IF EXISTS cities; +DROP SEQUENCE IF EXISTS cities_seq; DROP TABLE IF EXISTS users; DROP SEQUENCE IF EXISTS user_seq; DROP TYPE IF EXISTS user_flag; - CREATE TYPE user_flag AS ENUM ('active', 'deleted', 'superuser'); CREATE SEQUENCE user_seq START 100000; +CREATE SEQUENCE cities_seq START 100000; + +CREATE TABLE cities ( + id INTEGER PRIMARY KEY DEFAULT nextval('cities_seq'), + city TEXT NOT NULL +); +CREATE UNIQUE INDEX city_idx ON cities (city); CREATE TABLE users ( id INTEGER PRIMARY KEY DEFAULT nextval('user_seq'), full_name TEXT NOT NULL, email TEXT NOT NULL, - flag user_flag NOT NULL + flag user_flag NOT NULL, + city_id INTEGER NOT NULL, + CONSTRAINT city_fk FOREIGN KEY (city_id) REFERENCES cities (id) ); CREATE UNIQUE INDEX email_idx ON users (email); \ No newline at end of file diff --git a/persist/src/main/java/ru/javaops/masterjava/persist/dao/CityDao.java b/persist/src/main/java/ru/javaops/masterjava/persist/dao/CityDao.java new file mode 100644 index 0000000..4601376 --- /dev/null +++ b/persist/src/main/java/ru/javaops/masterjava/persist/dao/CityDao.java @@ -0,0 +1,42 @@ +package ru.javaops.masterjava.persist.dao; + +import com.bertoncelj.jdbi.entitymapper.EntityMapperFactory; +import org.skife.jdbi.v2.sqlobject.*; +import org.skife.jdbi.v2.sqlobject.customizers.BatchChunkSize; +import org.skife.jdbi.v2.sqlobject.customizers.RegisterMapperFactory; +import ru.javaops.masterjava.persist.DBIProvider; +import ru.javaops.masterjava.persist.model.City; + +import java.util.List; + +@RegisterMapperFactory(EntityMapperFactory.class) +public abstract class CityDao implements AbstractDao { + + public City insert(City city) { + if (city.isNew()) { + int id = insertGeneratedId(city); + city.setId(id); + } else { + insertWithId(city); + } + return city; + } + @SqlQuery("SELECT nextval('cities_seq')") + public abstract int getNextVal(); + + @SqlUpdate("TRUNCATE cities") + @Override + public abstract void clean(); + + @SqlUpdate("INSERT INTO cities (city) VALUES (:city) ON CONFLICT DO NOTHING") + @GetGeneratedKeys + abstract int insertGeneratedId(@BindBean City city); + + @SqlUpdate("INSERT INTO cities (id, city) VALUES (:id, :city) ON CONFLICT DO NOTHING") + @GetGeneratedKeys + abstract int insertWithId(@BindBean City city); + + @SqlQuery("SELECT * FROM cities LIMIT :it") + abstract List getWithLimit(@Bind int limit); + +} diff --git a/persist/src/main/java/ru/javaops/masterjava/persist/dao/UserDao.java b/persist/src/main/java/ru/javaops/masterjava/persist/dao/UserDao.java index 8932950..eb44116 100644 --- a/persist/src/main/java/ru/javaops/masterjava/persist/dao/UserDao.java +++ b/persist/src/main/java/ru/javaops/masterjava/persist/dao/UserDao.java @@ -26,7 +26,7 @@ public User insert(User user) { @SqlQuery("SELECT nextval('user_seq')") abstract int getNextVal(); - @Transaction + //@Transaction public int getSeqAndSkip(int step) { int id = getNextVal(); DBIProvider.getDBI().useHandle(h -> h.execute("ALTER SEQUENCE user_seq RESTART WITH " + (id + step))); @@ -49,7 +49,7 @@ public int getSeqAndSkip(int step) { public abstract void clean(); // https://habrahabr.ru/post/264281/ - @SqlBatch("INSERT INTO users (id, full_name, email, flag) VALUES (:id, :fullName, :email, CAST(:flag AS USER_FLAG))" + + @SqlBatch("INSERT INTO users (id, full_name, email, flag, city_id) VALUES (:id, :fullName, :email, CAST(:flag AS USER_FLAG), :city.id)" + "ON CONFLICT DO NOTHING") // "ON CONFLICT (email) DO UPDATE SET full_name=:fullName, flag=CAST(:flag AS USER_FLAG)") public abstract int[] insertBatch(@BindBean List users, @BatchChunkSize int chunkSize); diff --git a/persist/src/main/java/ru/javaops/masterjava/persist/model/City.java b/persist/src/main/java/ru/javaops/masterjava/persist/model/City.java new file mode 100644 index 0000000..6204942 --- /dev/null +++ b/persist/src/main/java/ru/javaops/masterjava/persist/model/City.java @@ -0,0 +1,16 @@ +package ru.javaops.masterjava.persist.model; + +import lombok.*; + +@Data +@EqualsAndHashCode(callSuper = true) +@RequiredArgsConstructor +@NoArgsConstructor +public class City extends BaseEntity { + private @NonNull String city; + + public City(Integer id, @NonNull String city) { + this(city); + this.id = id; + } +} diff --git a/persist/src/main/java/ru/javaops/masterjava/persist/model/Group.java b/persist/src/main/java/ru/javaops/masterjava/persist/model/Group.java new file mode 100644 index 0000000..8c63e88 --- /dev/null +++ b/persist/src/main/java/ru/javaops/masterjava/persist/model/Group.java @@ -0,0 +1,19 @@ +package ru.javaops.masterjava.persist.model; + +import com.bertoncelj.jdbi.entitymapper.Column; +import lombok.*; + +@Data +@RequiredArgsConstructor +@NoArgsConstructor +@EqualsAndHashCode(callSuper = true) +public class Group extends BaseEntity { + private @NonNull String name; + @Column("type_group") + private @NonNull TypeGroup typeGroup; + + public Group(Integer id, @NonNull String name, @NonNull TypeGroup typeGroup) { + this(name, typeGroup); + this.id = id; + } +} diff --git a/persist/src/main/java/ru/javaops/masterjava/persist/model/Project.java b/persist/src/main/java/ru/javaops/masterjava/persist/model/Project.java new file mode 100644 index 0000000..0b46a80 --- /dev/null +++ b/persist/src/main/java/ru/javaops/masterjava/persist/model/Project.java @@ -0,0 +1,17 @@ +package ru.javaops.masterjava.persist.model; + +import lombok.*; + +@Data +@RequiredArgsConstructor +@NoArgsConstructor +@EqualsAndHashCode(callSuper = true) +public class Project extends BaseEntity { + private @NonNull String name; + private @NonNull String description; + + public Project(Integer id, @NonNull String name, @NonNull String description) { + this(name, description); + this.id = id; + } +} diff --git a/persist/src/main/java/ru/javaops/masterjava/persist/model/TypeGroup.java b/persist/src/main/java/ru/javaops/masterjava/persist/model/TypeGroup.java new file mode 100644 index 0000000..ede17c9 --- /dev/null +++ b/persist/src/main/java/ru/javaops/masterjava/persist/model/TypeGroup.java @@ -0,0 +1,6 @@ +package ru.javaops.masterjava.persist.model; + +public enum TypeGroup { + FINISHED, + CURRENT +} diff --git a/persist/src/main/java/ru/javaops/masterjava/persist/model/User.java b/persist/src/main/java/ru/javaops/masterjava/persist/model/User.java index 927a299..0acb9ad 100644 --- a/persist/src/main/java/ru/javaops/masterjava/persist/model/User.java +++ b/persist/src/main/java/ru/javaops/masterjava/persist/model/User.java @@ -12,9 +12,9 @@ public class User extends BaseEntity { private @NonNull String fullName; private @NonNull String email; private @NonNull UserFlag flag; - - public User(Integer id, String fullName, String email, UserFlag flag) { - this(fullName, email, flag); + private @NonNull City city; + public User(Integer id, String fullName, String email, UserFlag flag, City city) { + this(fullName, email, flag, city); this.id=id; } } \ No newline at end of file diff --git a/persist/src/main/resources/persist.conf b/persist/src/main/resources/persist.conf index fe15029..a396e3d 100644 --- a/persist/src/main/resources/persist.conf +++ b/persist/src/main/resources/persist.conf @@ -4,4 +4,4 @@ db { password = password } -include required(file("/apps/masterjava/config/persist.conf")) +#include required(file("/apps/masterjava/config/persist.conf")) diff --git a/persist/src/test/java/ru/javaops/masterjava/persist/CityTestData.java b/persist/src/test/java/ru/javaops/masterjava/persist/CityTestData.java new file mode 100644 index 0000000..79f79b7 --- /dev/null +++ b/persist/src/test/java/ru/javaops/masterjava/persist/CityTestData.java @@ -0,0 +1,42 @@ +package ru.javaops.masterjava.persist; + +import com.google.common.collect.ImmutableList; +import ru.javaops.masterjava.persist.DBIProvider; +import ru.javaops.masterjava.persist.dao.CityDao; +import ru.javaops.masterjava.persist.model.City; + +import java.util.List; + +import static com.google.common.collect.ImmutableList.of; + +public class CityTestData { + public static City SPB; + public static City MSK; + public static City KIV; + public static City MNSK; + public static City Vladimir; + public static List fourCityList; + + public static void init() { + SPB = new City("Санкт-Петербург"); + MSK = new City("Москва"); + KIV = new City("Киев"); + MNSK = new City("Минск"); + Vladimir = new City("Владимир"); + fourCityList = ImmutableList.of(SPB, MSK, KIV, MNSK); + } + + public static void setUp() { + CityDao dao = DBIProvider.getDao(CityDao.class); + dao.clean(); + DBIProvider.getDBI().useTransaction((conn, status) -> { + fourCityList.forEach(dao::insert); + dao.insert(Vladimir); + }); + } + + public static City getNew() { + return new City("Екатеринбург"); + } + +} diff --git a/persist/src/test/java/ru/javaops/masterjava/persist/UserTestData.java b/persist/src/test/java/ru/javaops/masterjava/persist/UserTestData.java index 7c13c04..d16dfca 100644 --- a/persist/src/test/java/ru/javaops/masterjava/persist/UserTestData.java +++ b/persist/src/test/java/ru/javaops/masterjava/persist/UserTestData.java @@ -17,12 +17,12 @@ public class UserTestData { public static List FIST5_USERS; public static void init() { - ADMIN = new User("Admin", "admin@javaops.ru", UserFlag.superuser); - DELETED = new User("Deleted", "deleted@yandex.ru", UserFlag.deleted); - FULL_NAME = new User("Full Name", "gmail@gmail.com", UserFlag.active); - USER1 = new User("User1", "user1@gmail.com", UserFlag.active); - USER2 = new User("User2", "user2@yandex.ru", UserFlag.active); - USER3 = new User("User3", "user3@yandex.ru", UserFlag.active); + ADMIN = new User("Admin", "admin@javaops.ru", UserFlag.superuser, CityTestData.SPB); + DELETED = new User("Deleted", "deleted@yandex.ru", UserFlag.deleted, CityTestData.SPB); + FULL_NAME = new User("Full Name", "gmail@gmail.com", UserFlag.active, CityTestData.KIV); + USER1 = new User("User1", "user1@gmail.com", UserFlag.active, CityTestData.MSK); + USER2 = new User("User2", "user2@yandex.ru", UserFlag.active, CityTestData.KIV); + USER3 = new User("User3", "user3@yandex.ru", UserFlag.active, CityTestData.MNSK); FIST5_USERS = ImmutableList.of(ADMIN, DELETED, FULL_NAME, USER1, USER2); } diff --git a/persist/src/test/java/ru/javaops/masterjava/persist/dao/CityDaoTest.java b/persist/src/test/java/ru/javaops/masterjava/persist/dao/CityDaoTest.java new file mode 100644 index 0000000..4ae1707 --- /dev/null +++ b/persist/src/test/java/ru/javaops/masterjava/persist/dao/CityDaoTest.java @@ -0,0 +1,45 @@ +package ru.javaops.masterjava.persist.dao; + +import org.junit.Assert; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Test; +import ru.javaops.masterjava.persist.CityTestData; +import ru.javaops.masterjava.persist.model.City; + +import java.util.List; + +import static ru.javaops.masterjava.persist.CityTestData.getNew; + +public class CityDaoTest extends AbstractDaoTest { + + + public CityDaoTest() { + super(CityDao.class); + } + + @BeforeClass + public static void beforeClass() throws Exception { + CityTestData.init(); + } + + @Before + public void setUp() throws Exception { + CityTestData.setUp(); + } + + @Test + public void testInsert() { + City expected = getNew(); + City actual = dao.insert(expected); + Assert.assertNotNull(actual); + expected.setId(actual.getId()); + Assert.assertSame(actual, expected); + } + + @Test + public void testGetWithLimit() { + List actual = dao.getWithLimit(4); + Assert.assertEquals(CityTestData.fourCityList, actual); + } +} \ No newline at end of file diff --git a/web/upload/src/main/java/ru/javaops/masterjava/upload/CityProcessor.java b/web/upload/src/main/java/ru/javaops/masterjava/upload/CityProcessor.java new file mode 100644 index 0000000..3221d3f --- /dev/null +++ b/web/upload/src/main/java/ru/javaops/masterjava/upload/CityProcessor.java @@ -0,0 +1,58 @@ +package ru.javaops.masterjava.upload; + +import lombok.extern.slf4j.Slf4j; +import ru.javaops.masterjava.persist.DBIProvider; +import ru.javaops.masterjava.persist.dao.CityDao; +import ru.javaops.masterjava.persist.model.City; +import ru.javaops.masterjava.xml.schema.CityType; +import ru.javaops.masterjava.xml.schema.ObjectFactory; +import ru.javaops.masterjava.xml.util.JaxbParser; +import ru.javaops.masterjava.xml.util.JaxbUnmarshaller; +import ru.javaops.masterjava.xml.util.StaxStreamProcessor; + +import javax.xml.bind.JAXBException; +import javax.xml.stream.XMLStreamException; +import javax.xml.stream.events.XMLEvent; +import java.io.IOException; +import java.io.InputStream; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; + +@Slf4j +public class CityProcessor { + private static final int NUMBER_THREAD = 4; + private final ExecutorService executorService = Executors.newFixedThreadPool(NUMBER_THREAD); + private static final JaxbParser JAXB_PARSER = new JaxbParser(ObjectFactory.class); + private static final CityDao cityDao = DBIProvider.getDao(CityDao.class); + + public void process(final InputStream is) throws XMLStreamException, JAXBException, IOException { + log.info("Start processing parse cities"); + StaxStreamProcessor staxStreamProcessor = new StaxStreamProcessor(is); + JaxbUnmarshaller unmarshaller = JAXB_PARSER.createUnmarshaller(); + List> futuresCity = new ArrayList<>(); + while (staxStreamProcessor.doUntil(XMLEvent.START_ELEMENT, "City")) { + CityType cityType = unmarshaller.unmarshal(staxStreamProcessor.getReader(), CityType.class); + City city = new City(cityType.getValue()); + addCity(city, futuresCity); + } + List cities = new ArrayList<>(); + try { + for (Future f : futuresCity) { + cities.add(f.get()); + } + } catch (InterruptedException | ExecutionException e) { + + } + } + + private void addCity(City city, List> futuresCity) { + Future futureCity = executorService.submit(() -> cityDao.insert(city)); + futuresCity.add(futureCity); + log.info("city: {} insert", city.getCity()); + } + +} diff --git a/web/upload/src/main/java/ru/javaops/masterjava/upload/UploadServlet.java b/web/upload/src/main/java/ru/javaops/masterjava/upload/UploadServlet.java index 22c5f49..55fe719 100644 --- a/web/upload/src/main/java/ru/javaops/masterjava/upload/UploadServlet.java +++ b/web/upload/src/main/java/ru/javaops/masterjava/upload/UploadServlet.java @@ -23,7 +23,7 @@ public class UploadServlet extends HttpServlet { private static final int CHUNK_SIZE = 2000; private final UserProcessor userProcessor = new UserProcessor(); - + private final CityProcessor cityProcessor = new CityProcessor(); @Override protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws IOException { out(req, resp, "", CHUNK_SIZE); @@ -41,6 +41,7 @@ protected void doPost(HttpServletRequest req, HttpServletResponse resp) throws I } else { Part filePart = req.getPart("fileToUpload"); try (InputStream is = filePart.getInputStream()) { + cityProcessor.process(is); List failed = userProcessor.process(is, chunkSize); log.info("Failed users: " + failed); final WebContext webContext = diff --git a/web/upload/src/main/java/ru/javaops/masterjava/upload/UserProcessor.java b/web/upload/src/main/java/ru/javaops/masterjava/upload/UserProcessor.java index 93430af..a88bd96 100644 --- a/web/upload/src/main/java/ru/javaops/masterjava/upload/UserProcessor.java +++ b/web/upload/src/main/java/ru/javaops/masterjava/upload/UserProcessor.java @@ -5,6 +5,7 @@ import lombok.val; import ru.javaops.masterjava.persist.DBIProvider; import ru.javaops.masterjava.persist.dao.UserDao; +import ru.javaops.masterjava.persist.model.City; import ru.javaops.masterjava.persist.model.User; import ru.javaops.masterjava.persist.model.UserFlag; import ru.javaops.masterjava.xml.schema.ObjectFactory; @@ -14,6 +15,7 @@ import javax.xml.bind.JAXBException; import javax.xml.stream.XMLStreamException; import javax.xml.stream.events.XMLEvent; +import java.io.IOException; import java.io.InputStream; import java.util.ArrayList; import java.util.LinkedHashMap; @@ -45,9 +47,9 @@ public String toString() { /* * return failed users chunks */ - public List process(final InputStream is, int chunkSize) throws XMLStreamException, JAXBException { + public List process(final InputStream is, int chunkSize) throws XMLStreamException, JAXBException, IOException { log.info("Start processing with chunkSize=" + chunkSize); - + is.reset(); Map>> chunkFutures = new LinkedHashMap<>(); // ordered map (emailRange -> chunk future) int id = userDao.getSeqAndSkip(chunkSize); @@ -57,7 +59,8 @@ public List process(final InputStream is, int chunkSize) throws XM while (processor.doUntil(XMLEvent.START_ELEMENT, "User")) { ru.javaops.masterjava.xml.schema.User xmlUser = unmarshaller.unmarshal(processor.getReader(), ru.javaops.masterjava.xml.schema.User.class); - final User user = new User(id++, xmlUser.getValue(), xmlUser.getEmail(), UserFlag.valueOf(xmlUser.getFlag().value())); + City city = (City) xmlUser.getCity(); + final User user = new User(id++, xmlUser.getValue(), xmlUser.getEmail(), UserFlag.valueOf(xmlUser.getFlag().value()), (City) xmlUser.getCity()); chunk.add(user); if (chunk.size() == chunkSize) { addChunkFutures(chunkFutures, chunk); @@ -92,6 +95,5 @@ private void addChunkFutures(Map>> chunkFutures, Lis Future> future = executorService.submit(() -> userDao.insertAndGetConflictEmails(chunk)); chunkFutures.put(emailRange, future); log.info("Submit chunk: " + emailRange); - log.info("Submit chunk: " + emailRange); } }